跳至主要內容

createInngestAgent()

createInngestAgent() 為現有的 Agent 套用由 Inngest 驅動的持久執行包裝。它與 createDurableAgent() 一樣,透過 PubSub 串流傳送事件,並支援可恢復串流;但 Agent 迴圈會在 Inngest 的執行引擎上執行,而非在處理程序內執行。當一次執行必須不受處理程序重新啟動影響,或需要在分散式環境中執行時,請使用此函式。

如需處理程序內的持久執行,請使用 createDurableAgent()。如需在內置工作流程引擎上執行觸發後毋須等待結果的工作,請使用 createEventedAgent()

使用範例
使用範例 的直接連結

設定 Inngest 用戶端、包裝 Agent、向 Mastra 註冊,並公開 Inngest 服務端點:

src/mastra/index.ts
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:

Agent
要使用 Inngest 持久執行包裝的 Agent。InngestAgent 未有實作的方法(例如 listTools()getMemory())會透過 Proxy 委派給此 Agent。

inngest:

Inngest
Inngest 用戶端實例,用於傳送工作流程事件;在 SDK v4 中亦用於發佈即時串流事件。

id?:

string
= agent.id
覆寫 ID。

name?:

string
= agent.name
覆寫名稱。

pubsub?:

PubSub
= InngestPubSub
用於串流傳送事件的 PubSub 實例。預設的 InngestPubSub 使用 Inngest Realtime,可跨處理程序運作。

cache?:

MastraServerCache
用於儲存串流事件的快取,可啟用可恢復串流。提供此項後,PubSub 會自動以 CachingPubSub 包裝。如省略,Agent 會繼承 Mastra 實例的快取。

mastra?:

Mastra
用於可觀測性的 Mastra 實例。向 Mastra 註冊 Agent 時會自動設定。

InngestAgent 介面
inngestagent-interface 的直接連結

createInngestAgent() 傳回的物件,提供以下持久執行方法。任何未有明確定義的屬性或方法(例如 listTools()getMemory())都會透過 Proxy 轉送至底層 Agent。

屬性
屬性 的直接連結

id:

string
Agent ID。

name:

string
Agent 名稱。

agent:

Agent
底層 Mastra Agent。

inngest:

Inngest
Inngest 用戶端。

cache:

MastraServerCache | undefined
如已啟用可恢復串流,則為解析後的快取實例。

pubsub:

PubSub
用於串流傳送事件的 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()

第三個引數除了接受生命週期回調函式外,亦接受 threadIdresourceId

threadId?:

string
要與已恢復執行關聯的對話串 ID。

resourceId?:

string
要與已恢復執行關聯的資源 ID。

onChunk?:

(chunk: ChunkType) => void | Promise<void>
每個串流區塊傳送時呼叫。

onStepFinish?:

(result: AgentStepFinishEventData) => void | Promise<void>
agentic loop 中的步驟完成時呼叫。

onFinish?:

(result: AgentFinishEventData) => void | Promise<void>
執行完成時呼叫。

onError?:

(error: Error) => void | Promise<void>
執行發生錯誤時呼叫。

onSuspended?:

(data: AgentSuspendedEventData) => void | Promise<void>
執行暫停時呼叫。

傳回: Promise<InngestAgentStreamResult>

generate(messages, options?)
generatemessages-options 的直接連結

在 Inngest 的持久執行引擎上執行回應,並以單一 FullOutput 解析。如執行暫停,generate() 會以 finishReason: 'suspended' 解析。runId 選項並非必要;如省略,generate() 會建立執行 ID,並在 result.runId 中傳回。請使用 resumeGenerate() 繼續執行。如呼叫者需要暫停回調函式,請使用帶有 onSuspendedstream()

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() 的結果不包括 threadIdresourceId

傳回: 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?:

string
此次執行的唯一識別碼。稍後可配合 resume()observe() 使用。

instructions?:

AgentExecutionOptions['instructions']
覆寫 Agent 對此次執行的預設指示。

context?:

ModelMessage[]
提供給 Agent 的額外上下文訊息。

memory?:

object
用於持續保存及擷取對話的記憶體設定。

requestContext?:

RequestContext
帶有此次執行之動態設定及狀態的請求上下文。

maxSteps?:

number
最多可執行的步驟數目。

toolsets?:

object
此次執行可用的額外 Tool 集合。

clientTools?:

object
執行期間可用的用戶端 Tool。

toolChoice?:

'auto' | 'none' | 'required' | { type: 'tool'; toolName: string }
Tool 選擇策略。

modelSettings?:

object
模型特定設定,例如 temperature。

requireToolApproval?:

boolean
要求批准所有 Tool 呼叫;執行會暫停,直至恢復為止。

autoResumeSuspendedTools?:

boolean
自動恢復已暫停的 Tool,而非等候外部 resume() 呼叫。

toolCallConcurrency?:

number
可同時執行的 Tool 呼叫數目上限。

includeRawChunks?:

boolean
在串流輸出中包括原始 Provider 區塊。

maxProcessorRetries?:

number
每次生成中處理器重試次數上限。

untilIdle?:

boolean | { maxIdleMs?: number }
設定後,串流會在背景工作延續期間保持開啟,直至 Agent 閒置為止。傳入 true 可使用預設的 5 分鐘閒置逾時,或傳入 { maxIdleMs } 自訂。

onChunk?:

(chunk: ChunkType) => void | Promise<void>
每個串流區塊傳送時呼叫。

onStepFinish?:

(result: AgentStepFinishEventData) => void | Promise<void>
agentic loop 中的步驟完成時呼叫。

onFinish?:

(result: AgentFinishEventData) => void | Promise<void>
執行完成時呼叫。

onError?:

(error: Error) => void | Promise<void>
執行發生錯誤時呼叫。

onSuspended?:

(data: AgentSuspendedEventData) => void | Promise<void>
執行暫停時呼叫,例如等候 Tool 批准。

observe() 接受生命週期回調函式(onChunkonStepFinishonFinishonErroronSuspended),以及用於控制重播起點的 offset

InngestAgentStreamResult
inngestagentstreamresult 的直接連結

stream()resume() 傳回的物件。observe() 方法傳回相同結構,但不包括 threadIdresourceId

interface InngestAgentStreamResult<OUTPUT = undefined> {
output: MastraModelOutput<OUTPUT>
readonly fullStream: ReadableStream<any>
runId: string
threadId?: string
resourceId?: string
cleanup: () => void
}

output:

MastraModelOutput
串流輸出。等候 output.text 可取得完整文字,亦可取用 output.fullStream

fullStream:

ReadableStream
完整事件串流,委派至 output.fullStream

runId:

string
唯一執行 ID。將其傳給 resume()observe() 以重新連接。

threadId?:

string
使用記憶體時的對話串 ID。

resourceId?:

string
使用記憶體時的資源 ID。

cleanup:

() => void
取消訂閱 PubSub,並清除該次執行的登錄項目。完成該次執行後請呼叫此函式。

提供 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/expressinngest/fastifyinngest/next 等),並傳回該框架所用的服務函式。

Express
import { createServe } from '@mastra/inngest'
import { serve } from 'inngest/express'

const serveExpress = createServe(serve)
app.use('/inngest/api', serveExpress({ mastra, inngest }))
Next.js
import { createServe } from '@mastra/inngest'
import { serve } from 'inngest/next'

const serveNext = createServe(serve)
export const { GET, POST, PUT } = serveNext({ mastra, inngest })

服務選項
服務選項 的直接連結

mastra:

Mastra
包含已註冊 Agent 及 Workflow 的 Mastra 實例。

inngest:

Inngest
Inngest 用戶端實例。

functions?:

InngestFunction.Like[]
與 Mastra Workflow 一併提供的額外 Inngest 函式。

registerOptions?:

RegisterOptions
傳給 Inngest 註冊處理常式的選項。