Aller au contenu principal

createInngestAgent()

createInngestAgent() enveloppe un Agent existant avec une exécution durable fondée sur Inngest. Comme createDurableAgent(), il diffuse les événements via PubSub et prend en charge les streams reprenables, mais exécute la boucle agentique sur le moteur d'exécution d'Inngest plutôt que dans le processus. Utilisez-le lorsqu'une exécution doit survivre au redémarrage des processus ou s'exécuter dans un environnement distribué.

Pour une exécution durable dans le processus, utilisez createDurableAgent(). Pour une exécution sans attente sur le moteur de Workflow intégré, utilisez createEventedAgent().

Exemple d'utilisation
Lien direct vers Exemple d'utilisation

Configurez le client Inngest, enveloppez un Agent, enregistrez-le auprès de Mastra et exposez l'endpoint de service 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 }),
},
],
},
})

Diffusez une réponse et lisez le résultat :

const { output, runId, cleanup } = await durableAgent.stream('Hello!')

const text = await output.text

cleanup()

createInngestAgent(options)
Lien direct vers createinngestagentoptions

Enveloppe un Agent avec une exécution durable fondée sur Inngest et des streams reprenables.

import { createInngestAgent } from '@mastra/inngest'

const durableAgent = createInngestAgent({ agent, inngest })

Renvoie : InngestAgent

Paramètres
Lien direct vers Paramètres

agent:

Agent
Agent à envelopper avec l'exécution durable Inngest. Les méthodes non implémentées par InngestAgent, telles que listTools() et getMemory(), délèguent à cet Agent via un Proxy.

inngest:

Inngest
Instance du client Inngest. Utilisée pour envoyer les événements de Workflow et, dans le SDK v4, publier les événements du stream en temps réel.

id?:

string
= agent.id
Remplacement de l'identifiant.

name?:

string
= agent.name
Remplacement du nom.

pubsub?:

PubSub
= InngestPubSub
Instance PubSub destinée au streaming des événements. L'instance InngestPubSub par défaut utilise Inngest Realtime, qui fonctionne entre les processus.

cache?:

MastraServerCache
Cache des événements de stream stockés, qui permet de reprendre les streams. Lorsqu'il est fourni, PubSub est automatiquement enveloppé avec CachingPubSub. S'il est omis, l'Agent hérite du cache de l'instance Mastra.

mastra?:

Mastra
Instance Mastra destinée à l'Observability. Définie automatiquement lorsque l'Agent est enregistré auprès de Mastra.

Interface InngestAgent
Lien direct vers inngestagent-interface

Objet renvoyé par createInngestAgent(). Il fournit les méthodes d'exécution durable ci-dessous. Toute propriété ou méthode non définie explicitement, telle que listTools() et getMemory(), est transmise à l'Agent sous-jacent via un Proxy.

Propriétés
Lien direct vers Propriétés

id:

string
Identifiant de l'Agent.

name:

string
Nom de l'Agent.

agent:

Agent
Agent Mastra sous-jacent.

inngest:

Inngest
Client Inngest.

cache:

MastraServerCache | undefined
Instance de cache résolue, si les streams reprenables sont activés.

pubsub:

PubSub
Instance PubSub utilisée pour le streaming des événements.

Méthodes
Lien direct vers Méthodes

Exécution
Lien direct vers Exécution

stream(messages, options?)
Lien direct vers streammessages-options

Diffuse une réponse au moyen du moteur d'exécution durable d'Inngest. Le Workflow est déclenché via un événement Inngest après l'établissement de l'abonnement PubSub.

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()

Renvoie : Promise<InngestAgentStreamResult>

resume(runId, resumeData, options?)
Lien direct vers resumerunid-resumedata-options

Reprend une exécution Inngest suspendue, par exemple après l'approbation d'un Tool. Charge le snapshot du Workflow depuis le stockage, recherche l'étape suspendue et envoie un événement de reprise à Inngest.

const { output, cleanup } = await durableAgent.resume(
runId,
{
approved: true,
},
{ threadId: 'thread-1', resourceId: 'user-1' },
)

await output.text
cleanup()

Le troisième argument accepte threadId et resourceId en plus des callbacks du cycle de vie :

threadId?:

string
Identifiant du fil de discussion à associer à l'exécution reprise.

resourceId?:

string
Identifiant de la ressource à associer à l'exécution reprise.

onChunk?:

(chunk: ChunkType) => void | Promise<void>
Appelé pour chaque fragment diffusé.

onStepFinish?:

(result: AgentStepFinishEventData) => void | Promise<void>
Appelé lorsqu’une étape de la boucle agentique se termine.

onFinish?:

(result: AgentFinishEventData) => void | Promise<void>
Appelé lorsque l'exécution se termine.

onError?:

(error: Error) => void | Promise<void>
Appelé lorsque l'exécution rencontre une erreur.

onSuspended?:

(data: AgentSuspendedEventData) => void | Promise<void>
Appelé lorsque l'exécution est suspendue.

Renvoie : Promise<InngestAgentStreamResult>

generate(messages, options?)
Lien direct vers generatemessages-options

Exécute une réponse sur le moteur d'exécution durable d'Inngest et la résout en un unique FullOutput. Si l'exécution est suspendue, generate() est résolue avec finishReason: 'suspended'. L'option runId est facultative. Si vous l'omettez, generate() crée un identifiant d'exécution et le renvoie dans result.runId. Poursuivez l'exécution avec resumeGenerate(). Utilisez stream() avec onSuspended lorsque l'appelant a besoin d'un callback de suspension.

const result = await durableAgent.generate('Delete the old records', {
requireToolApproval: true,
})

result.runId // Generated automatically
result.finishReason // 'suspended' when approval is required

Renvoie : Promise<FullOutput<TOutput>>

resumeGenerate(runId, resumeData, options?)
Lien direct vers resumegeneraterunid-resumedata-options

Reprend une exécution generate() suspendue et la résout en un unique FullOutput.

if (!result.runId) {
throw new Error('Run ID is missing')
}

const resumedResult = await durableAgent.resumeGenerate(result.runId, { approved: true })

Renvoie : Promise<FullOutput<TOutput>>

observe(runId, options?)
Lien direct vers observerunid-options

Se reconnecte à une exécution existante, en rejouant les événements mis en cache avant de transmettre les événements en direct. Utilisez cette méthode après une déconnexion réseau. Transmettez offset pour commencer la relecture depuis une position connue.

const { output, cleanup } = await durableAgent.observe(runId, {
offset: 0,
onChunk: chunk => console.log(chunk),
})

await output.text

Le résultat de observe() n'inclut ni threadId ni resourceId.

Renvoie : Promise<Omit<InngestAgentStreamResult, 'threadId' | 'resourceId'>>

attention

Le cleanup() renvoyé par observe() détruit les entrées du registre et les événements en cache de l'exécution. Appelez-le uniquement lorsque vous avez terminé cette exécution. Si elle est suspendue et que vous comptez la reprendre plus tard, n'appelez pas cleanup().

prepare(messages, options?)
Lien direct vers preparemessages-options

Prépare une exécution durable sans la déclencher. Renvoie l'entrée sérialisée du Workflow qui permet de déclencher manuellement l'événement de Workflow Inngest.

const { runId, messageId, workflowInput, threadId, resourceId } = await durableAgent.prepare(
'Summarize the document',
{
memory: { threadId: 'thread-1', resourceId: 'user-1' },
},
)

Renvoie :

interface PrepareResult {
runId: string
messageId: string
workflowInput: any
threadId?: string
resourceId?: string
}

Introspection
Lien direct vers Introspection

isInngestAgent(obj)
Lien direct vers isinngestagentobj

Garde de type qui vérifie si un objet est un InngestAgent.

import { isInngestAgent } from '@mastra/inngest'

if (isInngestAgent(agent)) {
// agent is InngestAgent
}

Renvoie : boolean

Options du stream
Lien direct vers Options du stream

stream() accepte un objet InngestAgentStreamOptions. Il prend en charge les mêmes options d'exécution d'Agent que DurableAgent.stream(), ainsi que les callbacks du cycle de vie.

runId?:

string
Identifiant unique de cette exécution. Utilisez-le ensuite avec resume() ou observe().

instructions?:

AgentExecutionOptions['instructions']
Remplace les instructions par défaut de l'Agent pour cette exécution.

context?:

ModelMessage[]
Messages de contexte supplémentaires à fournir à l'Agent.

memory?:

object
Configuration de la mémoire pour la persistance et la récupération des conversations.

requestContext?:

RequestContext
Contexte de requête contenant la configuration dynamique et l'état de cette exécution.

maxSteps?:

number
Nombre maximal d'étapes à exécuter.

toolsets?:

object
Ensembles de Tools supplémentaires disponibles pour cette exécution.

clientTools?:

object
Tools côté client disponibles pendant l'exécution.

toolChoice?:

'auto' | 'none' | 'required' | { type: 'tool'; toolName: string }
Stratégie de sélection des Tools.

modelSettings?:

object
Paramètres propres au modèle, tels que la température.

requireToolApproval?:

boolean
Exige l'approbation de tous les appels de Tool, ce qui suspend l'exécution jusqu'à sa reprise.

autoResumeSuspendedTools?:

boolean
Reprend automatiquement les Tools suspendus au lieu d'attendre un appel externe à resume().

toolCallConcurrency?:

number
Nombre maximal d'appels de Tool à exécuter simultanément.

includeRawChunks?:

boolean
Inclut les fragments bruts du Provider dans la sortie du stream.

maxProcessorRetries?:

number
Nombre maximal de nouvelles tentatives des Processors par génération.

untilIdle?:

boolean | { maxIdleMs?: number }
Lorsque cette option est définie, maintient le stream ouvert pendant les continuations des tâches en arrière-plan jusqu'à ce que l'Agent soit inactif. Transmettez true pour le délai d'inactivité par défaut de 5 minutes, ou { maxIdleMs } pour le personnaliser.

onChunk?:

(chunk: ChunkType) => void | Promise<void>
Appelé pour chaque fragment diffusé.

onStepFinish?:

(result: AgentStepFinishEventData) => void | Promise<void>
Appelé lorsqu’une étape de la boucle agentique se termine.

onFinish?:

(result: AgentFinishEventData) => void | Promise<void>
Appelé lorsque l'exécution se termine.

onError?:

(error: Error) => void | Promise<void>
Appelé lorsque l'exécution rencontre une erreur.

onSuspended?:

(data: AgentSuspendedEventData) => void | Promise<void>
Appelé lorsque l'exécution est suspendue, par exemple pour l'approbation d'un Tool.

observe() accepte les callbacks du cycle de vie (onChunk, onStepFinish, onFinish, onError, onSuspended), ainsi qu'un offset contrôlant le point de départ de la relecture.

InngestAgentStreamResult
Lien direct vers inngestagentstreamresult

Objet renvoyé par stream() et resume(). La méthode observe() renvoie la même forme, mais omet threadId et resourceId.

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

output:

MastraModelOutput
Sortie en streaming. Attendez output.text pour obtenir le texte complet, ou consommez output.fullStream.

fullStream:

ReadableStream
Stream complet d'événements, qui délègue à output.fullStream.

runId:

string
Identifiant unique de l'exécution. Transmettez-le à resume() ou observe() pour vous reconnecter.

threadId?:

string
Identifiant du fil de discussion lors de l'utilisation de la mémoire.

resourceId?:

string
Identifiant de la ressource lors de l'utilisation de la mémoire.

cleanup:

() => void
Se désabonne de PubSub et efface les entrées du registre de l'exécution. Appelez cette fonction lorsque vous avez terminé l'exécution.

Mise à disposition des fonctions Inngest
Lien direct vers Mise à disposition des fonctions Inngest

Le package @mastra/inngest fournit serve() et createServe() pour enregistrer les fonctions de Workflow Inngest auprès de votre framework HTTP.

serve(options)
Lien direct vers serveoptions

Met à disposition les Workflows Mastra avec Hono, le framework par défaut. Collecte tous les Workflows Mastra fondés sur Inngest et les enregistre comme fonctions Inngest.

import { serve } from '@mastra/inngest'

app.use('/inngest/api', async c => {
return serve({ mastra, inngest })(c)
})

createServe(adapter)
Lien direct vers createserveadapter

Factory qui accepte n'importe quel adaptateur de service Inngest (inngest/express, inngest/fastify, inngest/next, etc.) et renvoie une fonction de service pour ce framework.

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 })

Options de service
Lien direct vers Options de service

mastra:

Mastra
Instance Mastra contenant les Agents et Workflows enregistrés.

inngest:

Inngest
Instance du client Inngest.

functions?:

InngestFunction.Like[]
Fonctions Inngest supplémentaires à mettre à disposition avec les Workflows Mastra.

registerOptions?:

RegisterOptions
Options transmises au gestionnaire d'enregistrement Inngest.