createInngestAgent()
createInngestAgent() 為現有的 Agent 套用由 Inngest 驅動的持久執行包裝。它與 createDurableAgent() 一樣,透過 PubSub 串流傳送事件,並支援可恢復串流;但 Agent 迴圈會在 Inngest 的執行引擎上執行,而非在處理程序內執行。當一次執行必須不受處理程序重新啟動影響,或需要在分散式環境中執行時,請使用此函式。
如需處理程序內的持久執行,請使用 createDurableAgent()。如需在內置工作流程引擎上執行觸發後毋須等待結果的工作,請使用 createEventedAgent()。
使用範例使用範例 的直接連結
設定 Inngest 用戶端、包裝 Agent、向 Mastra 註冊,並公開 Inngest 服務端點:
import { Mastra } from '@mastra/core'
import { Agent } from '@mastra/core/agent'
import { createInngestAgent, serve as inngestServe } from '@mastra/inngest'
import { Inngest } from 'inngest'
const inngest = new Inngest({ id: 'my-app' })
const agent = new Agent({
id: 'my-agent',
name: 'My Agent',
instructions: 'You are a helpful assistant',
model: 'openai/gpt-5.6-sol',
})
const durableAgent = createInngestAgent({ agent, inngest })
export const mastra = new Mastra({
agents: { myAgent: durableAgent },
server: {
apiRoutes: [
{
path: '/inngest/api',
method: 'ALL',
createHandler: async ({ mastra }) => inngestServe({ mastra, inngest }),
},
],
},
})
串流傳送回應並讀取結果:
const { output, runId, cleanup } = await durableAgent.stream('Hello!')
const text = await output.text
cleanup()
createInngestAgent(options)createinngestagentoptions 的直接連結
使用由 Inngest 驅動的持久執行及可恢復串流來包裝 Agent。
import { createInngestAgent } from '@mastra/inngest'
const durableAgent = createInngestAgent({ agent, inngest })
傳回: InngestAgent
參數參數 的直接連結
agent:
InngestAgent 未有實作的方法(例如 listTools() 和 getMemory())會透過 Proxy 委派給此 Agent。inngest:
id?:
name?:
pubsub?:
InngestPubSub 使用 Inngest Realtime,可跨處理程序運作。cache?:
CachingPubSub 包裝。如省略,Agent 會繼承 Mastra 實例的快取。mastra?:
InngestAgent 介面inngestagent-interface 的直接連結
createInngestAgent() 傳回的物件,提供以下持久執行方法。任何未有明確定義的屬性或方法(例如 listTools() 和 getMemory())都會透過 Proxy 轉送至底層 Agent。
屬性屬性 的直接連結
id:
name:
agent:
inngest:
cache:
pubsub:
方法方法 的直接連結
執行執行 的直接連結
stream(messages, options?)streammessages-options 的直接連結
使用 Inngest 的持久執行引擎串流傳送回應。建立 PubSub 訂閱後,工作流程會透過 Inngest 事件觸發。
const { output, runId, cleanup } = await durableAgent.stream('Hello!', {
onChunk: chunk => console.log(chunk),
onFinish: result => console.log('done', result),
})
const text = await output.text
cleanup()
傳回: Promise<InngestAgentStreamResult>
resume(runId, resumeData, options?)resumerunid-resumedata-options 的直接連結
恢復已暫停的 Inngest 執行,例如在 Tool 獲批後。此方法會從儲存空間載入工作流程快照、找出已暫停的步驟,並向 Inngest 傳送恢復事件。
const { output, cleanup } = await durableAgent.resume(
runId,
{
approved: true,
},
{ threadId: 'thread-1', resourceId: 'user-1' },
)
await output.text
cleanup()
第三個引數除了接受生命週期回調函式外,亦接受 threadId 和 resourceId:
threadId?:
resourceId?:
onChunk?:
onStepFinish?:
onFinish?:
onError?:
onSuspended?:
傳回: Promise<InngestAgentStreamResult>
generate(messages, options?)generatemessages-options 的直接連結
在 Inngest 的持久執行引擎上執行回應,並以單一 FullOutput 解析。如執行暫停,generate() 會以 finishReason: 'suspended' 解析。runId 選項並非必要;如省略,generate() 會建立執行 ID,並在 result.runId 中傳回。請使用 resumeGenerate() 繼續執行。如呼叫者需要暫停回調函式,請使用帶有 onSuspended 的 stream()。
const result = await durableAgent.generate('Delete the old records', {
requireToolApproval: true,
})
result.runId // Generated automatically
result.finishReason // 'suspended' when approval is required
傳回: Promise<FullOutput<TOutput>>
resumeGenerate(runId, resumeData, options?)resumegeneraterunid-resumedata-options 的直接連結
恢復已暫停的 generate() 執行,並以單一 FullOutput 解析。
if (!result.runId) {
throw new Error('Run ID is missing')
}
const resumedResult = await durableAgent.resumeGenerate(result.runId, { approved: true })
傳回: Promise<FullOutput<TOutput>>
observe(runId, options?)observerunid-options 的直接連結
重新連接現有執行,先重播快取事件,再傳送即時事件。網絡中斷後可使用此方法。傳入 offset 可從已知位置開始重播。
const { output, cleanup } = await durableAgent.observe(runId, {
offset: 0,
onChunk: chunk => console.log(chunk),
})
await output.text
observe() 的結果不包括 threadId 或 resourceId。
傳回: Promise<Omit<InngestAgentStreamResult, 'threadId' | 'resourceId'>>
observe() 傳回的 cleanup() 會刪除該次執行的登錄項目及快取事件。請只在完成該次執行後才呼叫。如執行已暫停,而你打算稍後恢復,請勿呼叫 cleanup()。
prepare(messages, options?)preparemessages-options 的直接連結
準備一次持久執行,但不會觸發執行。傳回已序列化的工作流程輸入,可用於手動觸發 Inngest 工作流程事件。
const { runId, messageId, workflowInput, threadId, resourceId } = await durableAgent.prepare(
'Summarize the document',
{
memory: { threadId: 'thread-1', resourceId: 'user-1' },
},
)
傳回:
interface PrepareResult {
runId: string
messageId: string
workflowInput: any
threadId?: string
resourceId?: string
}
內省內省 的直接連結
isInngestAgent(obj)isinngestagentobj 的直接連結
檢查物件是否為 InngestAgent 的型別守衛。
import { isInngestAgent } from '@mastra/inngest'
if (isInngestAgent(agent)) {
// agent is InngestAgent
}
傳回: boolean
串流選項串流選項 的直接連結
stream() 接受 InngestAgentStreamOptions 物件。除了生命週期回調函式外,它亦支援與 DurableAgent.stream() 相同的 Agent 執行選項。
runId?:
resume() 或 observe() 使用。instructions?:
context?:
memory?:
requestContext?:
maxSteps?:
toolsets?:
clientTools?:
toolChoice?:
modelSettings?:
requireToolApproval?:
autoResumeSuspendedTools?:
resume() 呼叫。toolCallConcurrency?:
includeRawChunks?:
maxProcessorRetries?:
untilIdle?:
true 可使用預設的 5 分鐘閒置逾時,或傳入 { maxIdleMs } 自訂。onChunk?:
onStepFinish?:
onFinish?:
onError?:
onSuspended?:
observe() 接受生命週期回調函式(onChunk、onStepFinish、onFinish、onError、onSuspended),以及用於控制重播起點的 offset。
InngestAgentStreamResultinngestagentstreamresult 的直接連結
stream() 和 resume() 傳回的物件。observe() 方法傳回相同結構,但不包括 threadId 和 resourceId。
interface InngestAgentStreamResult<OUTPUT = undefined> {
output: MastraModelOutput<OUTPUT>
readonly fullStream: ReadableStream<any>
runId: string
threadId?: string
resourceId?: string
cleanup: () => void
}
output:
output.text 可取得完整文字,亦可取用 output.fullStream。fullStream:
output.fullStream。runId:
resume() 或 observe() 以重新連接。threadId?:
resourceId?:
cleanup:
提供 Inngest 函式提供 Inngest 函式 的直接連結
@mastra/inngest 套件提供 serve() 和 createServe(),讓你在 HTTP 框架中註冊 Inngest 工作流程函式。
serve(options)serveoptions 的直接連結
使用 Hono(預設框架)提供 Mastra 工作流程。此函式會從 Mastra 收集所有由 Inngest 支援的工作流程,並將其註冊為 Inngest 函式。
import { serve } from '@mastra/inngest'
app.use('/inngest/api', async c => {
return serve({ mastra, inngest })(c)
})
createServe(adapter)createserveadapter 的直接連結
這個工廠函式接受任何 Inngest 服務適配器(inngest/express、inngest/fastify、inngest/next 等),並傳回該框架所用的服務函式。
import { createServe } from '@mastra/inngest'
import { serve } from 'inngest/express'
const serveExpress = createServe(serve)
app.use('/inngest/api', serveExpress({ mastra, inngest }))
import { createServe } from '@mastra/inngest'
import { serve } from 'inngest/next'
const serveNext = createServe(serve)
export const { GET, POST, PUT } = serveNext({ mastra, inngest })