createInngestAgent()
createInngestAgent()기존을 래핑Agent~와 함께섭취-강력한 내구성 실행. 좋다createDurableAgent(), 이벤트를 스트리밍합니다.PubSub재개 가능한 스트림을 지원하지만 프로세스 내 대신 Ingest의 실행 엔진에서 Agent 루프를 실행합니다. 실행이 프로세스를 다시 시작한 후에도 유지되어야 하거나 분산 환경에서 실행되어야 하는 경우 이를 사용합니다.
프로세스 내 내구성 실행에는 createDurableAgent()를 사용하세요. 기본 제공 Workflow 엔진에서 실행 후 결과를 기다리지 않는 방식으로 처리하려면 createEventedAgent()를 사용하세요.
사용예사용예에 대한 직접 링크
Ingest 클라이언트를 설정하고, Agent를 래핑하고, 이를 Mastra에 등록하고, Ingest 서비스 엔드포인트를 노출합니다.
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에 대한 직접 링크
Agent를 Inngest 기반 내구성 실행 및 재개 가능한 스트림으로 래핑합니다.
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에 대한 직접 링크
Ingest의 내구성 있는 실행 엔진을 사용하여 응답을 스트리밍합니다. Workflow는 PubSub 구독이 설정된 후 Ingest 이벤트를 통해 트리거됩니다.
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에 대한 직접 링크
예를 들어 Tool 승인 후 일시 중지된 Ingest 실행을 재개합니다. 스토리지에서 Workflow 스냅샷을 로드하고 일시 중단된 단계를 찾은 다음 재개 이벤트를 Ingest로 보냅니다.
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에 대한 직접 링크
트리거하지 않고 지속 가능한 실행을 위한 실행을 준비합니다. Ingest Workflow 이벤트를 수동으로 트리거하는 데 사용할 수 있는 직렬화된 Workflow 입력을 반환합니다.
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() 호출을 기다리지 않고 일시 중단된 Tool을 자동으로 재개합니다.toolCallConcurrency?:
includeRawChunks?:
maxProcessorRetries?:
untilIdle?:
true를 전달하고, 사용자 지정하려면 { 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:
Ingest 기능 제공Ingest 기능 제공에 대한 직접 링크
@mastra/inngest 패키지는 HTTP 프레임워크에 Inngest Workflow 함수를 등록하기 위한 serve()와 createServe()를 제공합니다.
serve(options)serveoptions에 대한 직접 링크
Hono(기본 프레임워크)를 사용하여 Mastra Workflow를 제공합니다. Mastra에서 모든 Ingest 지원 Workflow를 수집하고 이를 Ingest 기능으로 등록합니다.
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 등)를 받아 해당 프레임워크용 serve 함수를 반환하는 팩토리입니다.
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 })