Funkcjonalny interfejs API przepływu pracy

Warning

Funkcjonalny interfejs API przepływu pracy jest eksperymentalny i może ulec zmianie lub usunięciu w przyszłych wersjach bez powiadomienia.

Funkcjonalny interfejs API przepływu pracy umożliwia pisanie przepływów pracy jako zwykłych Python funkcji asynchronicznych. Zamiast definiować klasy wykonawcze, połączenia sterujące i korzystać z WorkflowBuilder, dekorujesz funkcję async za pomocą @workflow i używasz natywnego przepływu sterowania w Pythonie — if/else, pętli for, asyncio.gather — do wyrażenia swojej logiki.

Aby zapoznać się z porównaniem bezpośrednim z API grafu, zobacz Interfejsy API przepływu pracy w przeglądzie przepływu pracy.

@workflow Dekorator

Zastosuj @workflow do funkcji async, aby utworzyć bezstanowy FunctionalWorkflowDefinition:

from agent_framework import workflow

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

Dekorator @workflow obsługuje sparametryzowaną formę z opcjonalnymi argumentami:

from agent_framework import workflow

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

@workflow Parametry

Parametr Typ Opis
name str | None Nazwa przepływu pracy do wyświetlenia. Wartość domyślna funkcji to __name__.
description str | None Opcjonalny opis czytelny dla człowieka.

Tworzenie przepływu pracy

W definicji wywołaj metodę .build(), aby utworzyć komponent FunctionalWorkflow ze stanem. Każdy utworzony przepływ pracy reprezentuje jeden logiczny obiekt wywołujący lub sesję. Utwórz oddzielne wystąpienia dla niezależnych wywołujących, aby odizolować ich stan uruchamiania i odtwarzania.

Przekaż pamięć punktów kontrolnych do .build(), gdy przepływ pracy musi zachować wyniki i stan kroku:

from agent_framework import InMemoryCheckpointStorage

storage = InMemoryCheckpointStorage()
workflow_instance = text_pipeline.build(checkpoint_storage=storage)
Parametr Typ Opis
checkpoint_storage CheckpointStorage | None Domyślna pamięć masowa do przechowywania wyników kroków i stanu między uruchomieniami.

Sygnatura funkcji przepływu pracy

Pierwszy parametr funkcji przepływu pracy odbiera dane wejściowe przekazane do .run() utworzonego przepływu pracy. ctx: RunContext Dodaj parametr tylko wtedy, gdy potrzebujesz HITL, stanu klucz/wartość lub zdarzeń niestandardowych — w przeciwnym razie jest to opcjonalne:

# 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 Najpierw jest wykrywana przez adnotację typu, a następnie przez nazwę parametru ctx, więc zarówno ctx: RunContext, jak i samodzielny parametr ctx działają.

Uruchamianie przepływu pracy

Wywołaj .run() na obiekcie zwróconym przez FunctionalWorkflow.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() Parametry

Parametr Typ Opis
message Any | None Dane wejściowe przekazane do funkcji przepływu pracy jako pierwszy argument.
stream bool Jeśli True, zwraca ResponseStream, który generuje obiekty WorkflowEvent. Wartość domyślna to False.
responses dict[str, Any] | None Odpowiedzi HITL kluczowane przez request_id. Służy do wznawiania zawieszonego przepływu pracy.
checkpoint_id str | None Punkt kontrolny do przywrócenia. Wymaga ustawienia checkpoint_storage podczas kompilacji lub dla tego uruchomienia.
checkpoint_storage CheckpointStorage | None Nadpisuje pamięć punktów kontrolnych na etapie kompilacji dla tego uruchomienia.
include_status_events bool Uwzględnij zdarzenia zmiany stanu w wyniku braku przesyłania strumieniowego.

Podaj jeden tryb danych wejściowych dla każdego wywołania: message, responses lub checkpoint_id. Wyjątek to wznawianie punktu kontrolnego z danymi wejściowymi zewnętrznymi, gdzie checkpoint_id i responses można je przekazać razem.

WorkflowRunResult

run() (bez przesyłania strumieniowego) zwraca wartość WorkflowRunResult. Kluczowe metody:

Metoda / właściwość Returns Opis
.text str Pierwsze dane wyjściowe jako ciąg. Pusty ciąg, jeśli nie ma danych wyjściowych ciągu.
.get_outputs() list[Any] Wszystkie komunikaty wyjściowe terminala generowane przez przepływ pracy (zdarzenia zawierające type == "output").
.get_intermediate_outputs() list[Any] Wszystkie pośrednie dane wyjściowe generowane przez przepływ pracy (zdarzenia ze znacznikiem type == "intermediate").
.get_final_state() WorkflowRunState Stan ostatniego uruchomienia (IDLE, IDLE_WITH_PENDING_REQUESTS, FAILED, ...).
.get_request_info_events() list[WorkflowEvent] Oczekujące żądania HITL, gdy stan wynosi IDLE_WITH_PENDING_REQUESTS.

Streaming

Ustaw stream=True do odbierania zdarzeń w miarę ich tworzenia:

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

Zobacz python/samples/03-workflows/functional/basic_streaming_pipeline.py dla kompletnego przykładu.

@step Dekorator

@step to dekorator opcji, który dodaje możliwość buforowania wyników, emisję zdarzeń oraz punkty kontrolne dla poszczególnych kroków w funkcjach asynchronicznych:

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)

Co @step robi wewnątrz przepływu pracy

  • Buforuje wyniki — wynik jest przechowywany przez (step_name, call_index). Po wznowieniu lub przywróceniu punktu kontrolnego hitL ukończony krok zwraca zapisany wynik natychmiast zamiast ponownego wykonywania.
  • Emituje zdarzeniaexecutor_invoked / executor_completed / executor_failed są emitowane w celu obserwowania. W przypadku trafienia pamięci podręcznej, executor_bypassed jest emitowany zamiast tego.
  • Zapisuje punkty kontrolne — jeśli zbudowany przepływ pracy ma checkpoint_storage, punkt kontrolny jest zapisywany po ukończeniu każdego kroku.
  • Wstrzykuje RunContext — gdy funkcja kroku deklaruje parametr ctx: RunContext, aktywny kontekst jest automatycznie wstrzykiwany.

Poza uruchomionym procesem roboczym @step jest przezroczysty — funkcja zachowuje się identycznie jak jej nieozdobiona wersja, dzięki czemu jest w pełni testowana niezależnie.

Kiedy używać atrybutu @step

Użyj @step funkcji, które są kosztowne do powtórnego wykonania: wywołania agenta, żądania zewnętrznego interfejsu API lub dowolnej operacji, przy wznowieniu której ponowne wykonanie byłoby kosztowne lub miałoby skutki uboczne. Funkcje zwykłe (bez @step) nadal działają wewnątrz @workflow; po prostu są ponownie wykonywane po wznowieniu przepływu pracy.

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 również akceptuje parametr name:

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

Zobacz python/samples/03-workflows/functional/steps_and_checkpointing.py dla kompletnego przykładu.

RunContext

RunContext to kontekst wykonania wstrzykiwany do funkcji przepływu pracy i funkcji kroków. Potrzebujesz go tylko wtedy, gdy używasz funkcji HITL, stanu klucza/wartości lub zdarzeń niestandardowych.

Zaimportuj go z pliku agent_framework:

from agent_framework import RunContext, workflow

ctx.request_info() — pętla "human-in-the-loop"

ctx.request_info() wstrzymuje przepływ pracy do oczekiwania na dane wejściowe zewnętrzne:

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

Parametry:

Parametr Typ Opis
request_data Any Ładunek opisujący potrzebne dane wejściowe (dict, model Pydantic, string, ...).
response_type type Oczekiwany typ odpowiedzi w Pythonie.
request_id str | None Stabilny identyfikator dla tego żądania. W przypadku pominięcia identyfikator deterministyczny auto::<index> jest generowany na podstawie kolejności wywołań.

Semantyka odtwarzania: Przy pierwszym uruchomieniu request_info() zgłasza sygnał wewnętrzny, który zawiesza przepływ pracy. Sygnał nigdy nie jest widoczny dla Twojego kodu. Obiekt wywołujący otrzymuje element WorkflowRunResult z elementem get_final_state() == WorkflowRunState.IDLE_WITH_PENDING_REQUESTS.

Aby wznowić, wywołaj .run(responses={request_id: value}) dla tego samego zbudowanego przepływu pracy. Przepływ pracy jest wykonywany ponownie od góry i request_info() natychmiast zwraca podaną wartość.

@step-oznakowane funkcje, które działały przed zawieszeniem, zwracają buforowane wyniki podczas wznawiania zamiast ponownego wykonywania.

Obsługa odpowiedzi:

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)

Zobacz python/samples/03-workflows/functional/hitl_review.py dla kompletnego przykładu.

ctx.request_info() jest również obsługiwany wewnątrz funkcji @step.

ctx.add_event() — Zdarzenia niestandardowe

ctx.add_event() służy do emitowania zdarzeń specyficznych dla aplikacji równocześnie ze zdarzeniami cyklu życia platformy. Aby uzyskać szczegółowe informacje i przykłady, zobacz Emitowanie zdarzeń niestandardowych.

ctx.get_state() / ctx.set_state() — stan danych typu klucz/wartość

Użyj ctx.get_state() i ctx.set_state(), aby przechowywać wartości, które są zachowywane w przerwach HITL i w punktach kontrolnych. Aby uzyskać szczegółowe informacje, zobacz Stan przepływu pracy.

Wartości stanu muszą być serializowalne w formacie JSON po skonfigurowaniu magazynu punktów kontrolnych.

ctx.is_streaming()

Zwraca True, gdy bieżący przebieg został uruchomiony z stream=True. Przydatne wewnątrz funkcji kroków, które chcą dostosować swoje zachowanie na podstawie trybu przesyłania strumieniowego.

get_run_context()

Pobiera aktywne RunContext z dowolnego miejsca w uruchomionym przepływie pracy — przydatne w funkcjach pomocnika, które nie deklarują parametru ctx :

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)

Zwraca wartość None po wywołaniu poza działającym przepływem pracy.

Równoległość z asyncio.gather

Użyj standardowej współbieżności Python dla fan-out/fan-in — brak elementów pierwotnych platformy:

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 działa również wtedy, gdy funkcje są dekorowane @step.

Zobacz python/samples/03-workflows/functional/parallel_pipeline.py dla kompletnego przykładu.

Wywoływanie agentów wewnątrz przepływów pracy

Wywołania agenta działają jako zwykłe wywołania funkcji wewnątrz @workflowelementu :

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

Dodaj @step do funkcji wywołujących agentów, gdy ich wyniki mają być buforowane podczas wznowień HITL lub przy przywracaniu stanu.

from agent_framework import step

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

Zobacz python/samples/03-workflows/functional/agent_integration.py dla kompletnego przykładu.

.as_agent() — Używanie przepływu pracy jako agenta

Skompiluj obiekt FunctionalWorkflow, a następnie opakuj go jako obiekt zgodny z agentem za pomocą polecenia .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() Zwraca obiekt FunctionalWorkflowAgent , który uwidacznia ten sam run() interfejs co inne obiekty agenta, dzięki czemu funkcjonalne przepływy pracy są komponowalne z dowolnym systemem, który akceptuje agentów.

Parametr Typ Opis
name str | None Nazwa wyświetlana agenta. Domyślnie używana jest nazwa przepływu pracy.
description str | None Opcjonalne nadpisanie opisu. Domyślnie jest to opis przepływu pracy.
context_providers Sequence[Any] | None Opcjonalni dostawcy kontekstu do powiązania z agentem.

Samples

Przykłady z możliwością uruchamiania znajdują się w następujących folderach przykładowych:

Następne kroki

Powiązane tematy:

Funkcjonalny interfejs API przepływu pracy nie jest obecnie dostępny dla języka C#.

Funkcjonalny interfejs API przepływu pracy nie jest obecnie dostępny dla języka Go. Użyj interfejsów API grafowych przepływów pracy w Workflow Builder & Execution.