API del flusso di lavoro funzionale

Avvertimento

L'API del flusso di lavoro funzionale è sperimentale e soggetta a modifiche o rimozione nelle versioni future senza preavviso.

L'API del flusso di lavoro funzionale consente di scrivere flussi di lavoro come funzioni asincrone Python semplici. Invece di definire classi executor, configurare i collegamenti e usare WorkflowBuilder, si applica un decoratore a una funzione async con @workflow e si utilizza il flusso di controllo nativo di Python — if/else, cicli for, asyncio.gather — per esprimere la logica.

Per un confronto affiancato con l'API graph, vedere API del flusso di lavoro nella panoramica dei flussi di lavoro.

@workflow Decoratore

Applica @workflow a una funzione async per creare un oggetto FunctionalWorkflowDefinition senza stato:

from agent_framework import workflow

@workflow
async def text_pipeline(text: str) -> str:
    upper = await to_upper_case(text)
    return await reverse_text(upper)

L'elemento @workflow decorator supporta una forma parametrizzabile con argomenti facoltativi:

from agent_framework import workflow

@workflow(name="my_pipeline", description="Uppercase then reverse")
async def text_pipeline(text: str) -> str:
    ...

@workflow Parametri

Parametro Type Descrizione
name str | None Nome visualizzato per il flusso di lavoro. Per impostazione predefinita, viene usato l'elemento __name__ della funzione.
description str | None Descrizione leggibile facoltativa.

Creare un flusso di lavoro

Chiama .build() sulla definizione per creare un FunctionalWorkflow con stato. Ogni flusso di lavoro compilato rappresenta un chiamante logico o una sessione. Crea istanze separate per i chiamanti indipendenti, in modo da isolare il loro stato di esecuzione e di replay.

Passa l'archiviazione dei checkpoint a .build() quando il flusso di lavoro deve salvare in modo persistente i risultati dei passaggi e lo stato:

from agent_framework import InMemoryCheckpointStorage

storage = InMemoryCheckpointStorage()
workflow_instance = text_pipeline.build(checkpoint_storage=storage)
Parametro Type Descrizione
checkpoint_storage CheckpointStorage | None Archiviazione predefinita per rendere persistenti i risultati dei passaggi e lo stato tra le esecuzioni.

Firma della funzione flusso di lavoro

Il primo parametro della funzione del flusso di lavoro riceve l'input passato a .run() nel flusso di lavoro compilato. Aggiungere un ctx: RunContext parametro solo quando è necessario HITL, lo stato key/value o eventi personalizzati — altrimenti è facoltativo:

# No ctx needed — just a plain pipeline
@workflow
async def simple_pipeline(data: str) -> str:
    result = await process(data)
    return result

# ctx needed for HITL, state, or custom events
@workflow
async def hitl_pipeline(data: str, ctx: RunContext) -> str:
    feedback = await ctx.request_info({"draft": data}, response_type=str)
    return feedback

RunContext viene rilevato prima dall'annotazione di tipo, quindi dal nome del parametro ctx, quindi funzionano sia ctx: RunContext che un parametro semplice ctx.

Esecuzione di un flusso di lavoro

Chiamare .run() sull'oggetto FunctionalWorkflow restituito da .build():

workflow_instance = text_pipeline.build()

# .run() wraps the result in a WorkflowRunResult with events and state
result = await workflow_instance.run("hello world")
print(result.text)                        # first output as a string
print(result.get_outputs())               # list of terminal outputs
print(result.get_intermediate_outputs())  # list of intermediate outputs
print(result.get_final_state())           # WorkflowRunState.IDLE

run() Parametri

Parametro Type Descrizione
message Any | None Input passato alla funzione del workflow come primo argomento.
stream bool Se True, fornisce un ResponseStream che restituisce oggetti WorkflowEvent. Di default è False.
responses dict[str, Any] | None Risposte HITL indicizzate da request_id. Utilizzato per riprendere un flusso di lavoro sospeso.
checkpoint_id str | None Checkpoint da cui eseguire il ripristino. Richiede che checkpoint_storage sia impostato durante la compilazione o per questa esecuzione.
checkpoint_storage CheckpointStorage | None Sostituisce la memorizzazione del checkpoint in fase di compilazione per questa esecuzione.
include_status_events bool Includere eventi di modifica dello stato nel risultato non di streaming.

Specificare una modalità di input per chiamata: message, responseso checkpoint_id. L'eccezione è la ripresa da checkpoint con input esterno, in cui checkpoint_id e responses possono essere passati insieme.

WorkflowRunResult

run() (senza streaming) restituisce un elemento WorkflowRunResult. Metodi chiave:

Metodo/proprietà Restituzioni Descrizione
.text str Primo output come stringa. Stringa vuota se non sono presenti output stringa.
.get_outputs() list[Any] Tutti gli output del terminale generati dal flusso di lavoro (eventi con type == "output").
.get_intermediate_outputs() list[Any] Tutti gli output intermedi generati dal flusso di lavoro (eventi con type == "intermediate").
.get_final_state() WorkflowRunState Stato di esecuzione finale (IDLE, IDLE_WITH_PENDING_REQUESTS, FAILED, ...).
.get_request_info_events() list[WorkflowEvent] Richieste HITL in sospeso quando lo stato è IDLE_WITH_PENDING_REQUESTS.

Streaming

Passare stream=True per ricevere gli eventi man mano che vengono generati:

from agent_framework import workflow

@workflow
async def data_pipeline(url: str) -> str:
    raw = await fetch_data(url)
    return await transform_data(raw)

workflow_instance = data_pipeline.build()

# stream=True returns a ResponseStream you iterate with async for
stream = workflow_instance.run("https://example.com/api/data", stream=True)
async for event in stream:
    if event.type == "output":
        print(f"Output: {event.data}")

# After iteration, get_final_response() returns the WorkflowRunResult
result = await stream.get_final_response()
print(f"Final state: {result.get_final_state()}")

Vedere python/samples/03-workflows/functional/basic_streaming_pipeline.py per un esempio completo.

@step Decoratore

@step è un decoratore da attivare esplicitamente che aggiunge la memorizzazione nella cache dei risultati, l'emissione di eventi e la creazione di checkpoint per passaggio alle singole funzioni asincrone:

from agent_framework import step, workflow

@step
async def fetch_data(url: str) -> dict:
    # expensive — hits a real API
    return await http_get(url)

@workflow
async def pipeline(url: str) -> str:
    raw = await fetch_data(url)
    return process(raw)

Operazioni @step all'interno di un flusso di lavoro

  • Memorizza nella cache i risultati : il risultato viene archiviato da (step_name, call_index). In caso di riprendimento o ripristino del checkpoint HITL, un passaggio completato restituisce immediatamente il risultato salvato anziché eseguire di nuovo.
  • Genera eventi : executor_invoked / executor_completed / executor_failed vengono generati per l'osservabilità. In caso di hit della cache, viene invece emesso executor_bypassed.
  • Salva i checkpoint : se il flusso di lavoro compilato include checkpoint_storage, viene salvato un checkpoint al termine di ogni passaggio.
  • Inietta RunContext — se la funzione step dichiara un ctx: RunContext parametro, il contesto attivo viene inserito automaticamente.

All'esterno di un flusso di lavoro in esecuzione, @step è trasparente: la funzione si comporta in modo identico alla versione non decorata, rendendola completamente testabile in isolamento.

Quando usare @step

Usare @step sulle funzioni che sono costose da rieseguire: chiamate agli agenti, richieste API esterne o qualsiasi operazione in cui la riesecuzione alla ripresa sarebbe costosa o avrebbe effetti collaterali. Le funzioni semplici (senza @step) funzionano ancora all'interno @workflowdi ; vengono semplicemente eseguite nuovamente quando il flusso di lavoro riprende.

from agent_framework import InMemoryCheckpointStorage, step, workflow

@step  # cached — won't re-run on resume
async def call_llm(prompt: str) -> str:
    return (await agent.run(prompt)).text

# No @step — cheap, fine to re-run
async def validate(text: str) -> bool:
    return len(text) > 0

@workflow
async def pipeline(topic: str) -> str:
    draft = await call_llm(f"Write about: {topic}")
    ok = await validate(draft)
    return draft if ok else ""

storage = InMemoryCheckpointStorage()
workflow_instance = pipeline.build(checkpoint_storage=storage)

@step accetta anche un name parametro:

@step(name="transform")
async def transform_data(raw: dict) -> str:
    ...

Vedere python/samples/03-workflows/functional/steps_and_checkpointing.py per un esempio completo.

RunContext

RunContext è il contesto di esecuzione iniettato nelle funzioni del workflow e dello step. È necessario solo quando si usa HITL, stato chiave-valore o eventi personalizzati.

Importarlo da agent_framework:

from agent_framework import RunContext, workflow

ctx.request_info() — Interazione umana nel processo

ctx.request_info() sospende il flusso di lavoro per attendere l'input esterno:

@workflow
async def review_pipeline(topic: str, ctx: RunContext) -> str:
    draft = await write_draft(topic)
    feedback = await ctx.request_info(
        {"draft": draft, "instructions": "Please review this draft"},
        response_type=str,
        request_id="review_request",
    )
    return await revise_draft(draft, feedback)

Parametri:

Parametro Type Descrizione
request_data Any Payload che descrive l'input necessario (dict, modello Pydantic, stringa, ...).
response_type type Tipo Python previsto della risposta.
request_id str | None Identificatore stabile per questa richiesta. Se omesso, viene generato un ID deterministico auto::<index> dall'ordine di chiamata.

Semantica di riesecuzione: Alla prima esecuzione, request_info() emette un segnale interno che sospende il workflow. Il segnale non è mai visibile dal tuo codice. Il chiamante riceve un WorkflowRunResult con get_final_state() == WorkflowRunState.IDLE_WITH_PENDING_REQUESTS.

Per riprendere, richiamare .run(responses={request_id: value}) nello stesso workflow compilato. Il flusso di lavoro viene eseguito nuovamente dall'inizio e request_info() restituisce immediatamente il valore specificato.

Le funzioni decorate con @step, che sono state eseguite prima della sospensione, restituiscono i loro risultati memorizzati nella cache al riavvio invece di essere rieseguite.

Gestione della risposta:

workflow_instance = review_pipeline.build()

# Phase 1 — run until the workflow pauses
result1 = await workflow_instance.run("AI Safety")
assert result1.get_final_state() == WorkflowRunState.IDLE_WITH_PENDING_REQUESTS

requests = result1.get_request_info_events()
print(requests[0].request_id)  # "review_request"

# Phase 2 — resume with the human's answer
result2 = await workflow_instance.run(
    responses={"review_request": "Add more details about alignment research"}
)
print(result2.text)

Vedere python/samples/03-workflows/functional/hitl_review.py per un esempio completo.

ctx.request_info() è supportato anche all'interno delle funzioni @step.

ctx.add_event() — Eventi personalizzati

Usare ctx.add_event() per generare eventi specifici dell'applicazione insieme agli eventi del ciclo di vita del framework. Per informazioni dettagliate ed esempi, vedere Creazione di eventi personalizzati.

ctx.get_state() / ctx.set_state() — Stato chiave/valore

Usare ctx.get_state() e ctx.set_state() per archiviare i valori che vengono mantenuti tra le interruzioni HITL e sono inclusi nei checkpoint. Per informazioni dettagliate, vedere Stato del flusso di lavoro.

I valori di stato devono essere serializzabili in JSON quando è configurata l'archiviazione del checkpoint.

ctx.is_streaming()

Restituisce True quando l'esecuzione corrente è stata avviata con stream=True. Utile all'interno delle funzioni dei passaggi che devono modificare il proprio comportamento in base alla modalità di streaming.

get_run_context()

Recupera l'oggetto attivo RunContext da qualsiasi posizione all'interno di un flusso di lavoro in esecuzione, utile nelle funzioni helper che non dichiarano un ctx parametro:

from agent_framework import get_run_context

async def helper():
    ctx = get_run_context()
    if ctx is not None:
        ctx.set_state("helper_ran", True)

Restituisce None quando viene chiamato all'esterno di un flusso di lavoro in esecuzione.

Parallelismo con asyncio.gather

Usare la concorrenza Python standard per fan-out/fan-in: non sono necessarie primitive del framework:

import asyncio
from agent_framework import workflow

@workflow
async def research_pipeline(topic: str) -> str:
    web, papers, news = await asyncio.gather(
        research_web(topic),
        research_papers(topic),
        research_news(topic),
    )
    return await synthesize([web, papers, news])

asyncio.gather funziona anche quando le funzioni sono decorate con @step.

Vedere python/samples/03-workflows/functional/parallel_pipeline.py per un esempio completo.

Chiamata di agenti all'interno dei flussi di lavoro

Le chiamate di Agent funzionano come chiamate di funzione normali all'interno di @workflow.

from agent_framework import Agent, workflow

writer = Agent(name="WriterAgent", instructions="Write a short poem.", client=client)
reviewer = Agent(name="ReviewerAgent", instructions="Review the poem.", client=client)

@workflow
async def poem_workflow(topic: str) -> str:
    poem = (await writer.run(f"Write a poem about: {topic}")).text
    review = (await reviewer.run(f"Review this poem: {poem}")).text
    return f"Poem:\n{poem}\n\nReview: {review}"

Aggiungere @step alle funzioni che chiamano agenti quando si vuole che i risultati vengano memorizzati nella cache tra riprese HITL o ripristini da checkpoint:

from agent_framework import step

@step
async def write_poem(topic: str) -> str:
    return (await writer.run(f"Write a poem about: {topic}")).text

Vedere python/samples/03-workflows/functional/agent_integration.py per un esempio completo.

.as_agent() — Uso di un flusso di lavoro come agente

Crea un FunctionalWorkflow, quindi racchiudilo in un oggetto compatibile con l'agente con .as_agent():

from agent_framework import workflow

@workflow
async def poem_workflow(topic: str) -> str:
    ...

# Build the workflow and wrap it as an agent
agent = poem_workflow.build().as_agent(name="PoemAgent")

# Use with the standard agent interface
response = await agent.run("Write a poem about the ocean")
print(response.text)

# Or use in a larger workflow or orchestration

.as_agent() restituisce un FunctionalWorkflowAgent oggetto che espone la stessa run() interfaccia di altri oggetti agente, rendendo componibili flussi di lavoro funzionali con qualsiasi sistema che accetta gli agenti.

Parametro Type Descrizione
name str | None Mostra il nome dell'agente. Per impostazione predefinita, viene usato il nome del workflow.
description str | None Sovrascrittura della descrizione facoltativa. Per impostazione predefinita, viene utilizzata la descrizione del flusso di lavoro.
context_providers Sequence[Any] | None Provider di contesto opzionali da associare all'agente.

Samples

Gli esempi eseguibili sono disponibili nelle cartelle di esempio seguenti:

Passaggi successivi

Argomenti correlati:

L'API del flusso di lavoro funzionale non è attualmente disponibile per C#.

L'API del flusso di lavoro funzionale non è attualmente disponibile per Go. Usa le API del flusso di lavoro grafico in Creazione ed esecuzione dei flussi di lavoro.