Skip to main content

LDF — LLM Dispatching Framework

An essential async LLM task dispatching and orchestration library.

This library is designed to be as simple as possible while being used in programming workflows that require both dynamic and deterministic execution paths.

It provides simple and intuitive usable classes for orchestrating tasks where an LLM is required (e.g. for explainability, data formatting, summarization, routing, etc.).

The idea behind this library is to use and implement LLM calls only when needed (no agentic loops or tool calling). This framework allows for flexible task orchestration, while keeping token usage predictable and efficient.

It has only one dependency: openai.

Installation

pip install -e .

Quick Start

import asyncio
from ldf import Task

async def main():
    task = Task("geography", "What is the capital of France?", system_prompt="You are a geography expert.")
    await task.run()
    print(task.output)

asyncio.run(main())

Streaming

async def on_chunk(chunk: str):
    print(chunk, end="", flush=True)

task = Task("story", "Tell me a short story", streaming=True, on_chunk=on_chunk)
await task.run()

If on_chunk is not provided, this default print behavior is used.

Iterating on a task (ReAct pattern idea)

n_iter feeds each iteration's output back into the LLM as the next prompt:

task = Task("story", "Write one sentence of a story.", n_iter=3)
await task.run()
for step in task.iterations:
    print(step["iteration"], step["output"])

run() resets iterations and output at the start of each call, so it's safe to await the same Task/Function instance again later (e.g. reusing module-level task objects across multiple requests) — state from a previous run never leaks into the next one.

Grouping independent tasks with MultiTask (parallel execution)

MultiTask runs its tasks concurrently via asyncio.gather. Only use it for tasks that don't depend on each other's output — since they run at the same time, none of them can rely on another's result being in context yet:

from ldf import Task, MultiTask

group = MultiTask([
    Task("summary", "Summarize: ..."),
    Task("translation", "Translate to French: ..."),
])

await group.run()
print(group)

If you need a strict dependency between tasks (one reads another's output via context), don't use MultiTask — just list them in order in a Workflow's tasks, which always runs sequentially.

Composing a Workflow

from ldf import Task, MultiTask, Workflow

workflow = Workflow("my_workflow", [
    Task("step_1", "Step 1"),
    MultiTask([Task("step_2a", "Step 2a"), Task("step_2b", "Step 2b")]),
])

await workflow.run()

Workflows can also nest other workflows, since a Workflow is itself a valid task in a tasks list:

sub_workflow = Workflow("sub_workflow", [Task("step_1", "Step 1")])

workflow = Workflow("main_workflow", [
    sub_workflow,
    Task("step_2", "Step 2"),
])

await workflow.run()

Conditions and context

Every Task, MultiTask, Workflow, and Function runs against a shared context dict, which is passed down through run() and populated with each named task's result (context[task.name] = task) as it completes.

A condition callable receives that context and decides whether the task should run. It can inspect the output of any earlier task by name:

route_task = Task("route", "Classify this request: chat or play_movie")

chat_task = Task(
    "chat", "Reply to the user",
    condition=lambda ctx: ctx["route"].output == "chat",
)

If condition returns a falsy value, the task is skipped (but still registered in the context under its name). condition may also be an async callable.

Passing one task's result into another

user_input accepts a callable (sync or async) instead of a plain string. It's resolved lazily against context right before the first iteration runs — so it can reference the output of any task that ran earlier in the same Workflow's tasks list:

route_task = Task("route", "Classify this request: chat or play_movie")

extract_task = Task(
    "extract_movie", "Extract the movie title mentioned in this request",
    condition=lambda ctx: ctx["route"].output == "play_movie",
)

summary_task = Task(
    "summary",
    user_input=lambda ctx: f"Write a one-line blurb for the movie: {ctx['extract_movie'].output}",
    condition=lambda ctx: ctx["extract_movie"].output,
)

This mirrors condition's pattern exactly (a callable taking context, awaited if it returns an awaitable), and only applies to the first iteration's prompt — later n_iter iterations still feed off the previous iteration's output. The resolved string is recorded in task.iterations[0]["input"] like any other prompt.

Accessing the original input from anywhere

Pass user_input to Workflow itself and it's automatically seeded into context["input"] before any of its tasks run — so every condition or dynamic user_input callable in the workflow can read it, with no need to thread it through manually via run(context=...):

workflow = Workflow(
    "chatbot_workflow",
    tasks=[route_task, chat_task, extract_task, play_task],
    user_input=user_input,
)

await workflow.run()
chat_task = Task(
    "chat", "...",
    condition=lambda ctx: ctx["route"].output == "chat" and ctx["input"],
)

Running plain functions with Function class

Function wraps a regular (sync or async) callable so it can sit inside a Workflow alongside Tasks, with the same condition/context support:

from ldf import Function

def play_movie(ctx):
    movie = ctx["extract_movie"].output
    print("Playing " + movie + "...")

play_movie_function = Function(
    name="play_movie",
    func=play_movie,
    condition=lambda ctx: ctx["extract_movie"].output,
)

Tools: LLM-assisted function execution

Tool wraps a deterministic function with an LLM-based parameter extraction and result analysis loop. Given a user request, the Tool:

  1. Uses an LLM Task to extract the function's parameters from the user input
  2. Executes the deterministic function with those parameters via a Function
  3. Uses an LLM Task to analyze whether the result is valid and useful

If the result is invalid, it retries up to max_attempts times.

from ldf import Tool

def extract_ips(text: str) -> str:
    import re
    return str(re.findall(r'\d+\.\d+\.\d+\.\d+', text))

tool = Tool(
    name="ip_extractor",
    description="Extracts all IP addresses from a text. Pass the text to search.",
    function=extract_ips,
    max_attempts=3,
    streaming=True,
)

output = await tool.run("Find IPs in: 192.168.1.1 and 10.0.0.5")
print(output)

Optional validators can clean up LLM-extracted parameters or function output:

tool = Tool(
    name="ip_classifier",
    description="Classifies an IP address. Pass a single IP.",
    function=classify_ip,
    parameter_validator=lambda p: p.strip().strip('"'),
    output_validator=lambda o: o.lower(),
)

Deterministic parameter extraction and result analysis

For most tools, parameters can be extracted deterministically (e.g. with regex) instead of using an LLM call. The parameter_extractor and result_analyzer callables replace the LLM-based tasks entirely — no tokens consumed, no non-determinism:

import re

tool = Tool(
    name="ip_extractor",
    description="Extracts all IP addresses from a text.",
    function=extract_ips,
    parameter_extractor=lambda user_input: user_input,  # pass input directly
    result_analyzer=lambda output: bool(output and output != "[]"),
)

When parameter_extractor is provided, the Tool calls it directly instead of running the LLM-based extract_parameters_task. When result_analyzer is provided, it replaces the LLM-based analyze_results_task. Both default to None, falling back to LLM — so you can mix deterministic and LLM-based tools in the same agent.

Skipping result analysis entirely

For tools where the function output is always valid (e.g. deterministic functions like regex extraction, hashing, arithmetic), set skip_analysis=True to return the result immediately after execution — no LLM call, no result_analyzer, no retry loop:

tool = Tool(
    name="sha256",
    description="Computes the SHA256 hash of a text.",
    function=compute_hash,
    parameter_extractor=lambda user_input: user_input,
    skip_analysis=True,
)

This saves one LLM call per tool execution. The tool still runs parameter extraction and the function itself, then returns the output directly.

Agents: LLM routing + Tool execution

Agent orchestrates multiple Tools with an LLM-based routing loop:

  1. An LLM Task routes the user input to one or more tools (by name, comma-separated)
  2. Each routed tool executes (using its own parameter extraction + execution + analysis loop)
  3. An LLM Task analyzes whether the combined results fulfill the user's request
  4. If not, it re-routes and re-executes up to max_attempts times
from ldf import Agent, Tool

agent = Agent(
    name="forensic_agent",
    user_input="Perform a forensic audit on this log: ...",
    tools=[
        Tool(name="ip_extractor", description="Extracts IPs from text", function=extract_ips,
             parameter_extractor=lambda text: text, result_analyzer=lambda o: bool(o)),
        Tool(name="email_extractor", description="Extracts emails from text", function=extract_emails,
             parameter_extractor=lambda text: text, result_analyzer=lambda o: bool(o)),
    ],
    max_attempts=5,
    streaming=True,
)

result = await agent.run()
print(result)

Optional route_validator and result_validator callables can clean up or filter the LLM's routing output and final result:

valid_tools = {"ip_extractor", "email_extractor"}

agent = Agent(
    name="forensic_agent",
    user_input=user_request,
    tools=tools,
    route_validator=lambda route: ", ".join(
        t for t in route.split(",") if t.strip() in valid_tools
    ),
)

Using MCP servers as tool sources

Agent also accepts an MCP instance directly as the tools argument. The agent automatically calls parse() if it hasn't been called yet, and uses the discovered tools:

from ldf import Agent, MCP

mcp = MCP("path/to/mcp.json")
agent = Agent(
    name="my_agent",
    user_input="Read /tmp/example.txt and search for errors",
    tools=mcp,
    max_attempts=5,
    streaming=True,
)

result = await agent.run()
print(result)

The agent automatically closes MCP connections when run() completes. See the MCP section below for details on configuring MCP servers.

Retries

Both Task and Function accept max_retries (default 0, no retry) and retry_delay (default 1.0 seconds). On failure they retry with exponential backoff (retry_delay * 2 ** attempt) before raising:

task = Task(
    "search", "Look this up...",
    max_retries=3,
    retry_delay=1.0,  # retries after ~1s, 2s, 4s
)

Retries apply per LLM/function call, not to the whole n_iter loop — if iteration 2 of 3 fails, only that call is retried, not iterations 0-1. Once max_retries is exhausted, the original exception propagates as before.

Logging

Task, Function, and LLM log failures and retry attempts through the standard logging module, under logger names ldf.task, ldf.function, and ldf.llm. LDF never configures handlers or exporters itself — that's left entirely to the consuming application, so you can wire it into whatever observability stack you use.

To bridge these logs into OpenTelemetry, attach the OTel SDK's LoggingHandler to the root logger (requires opentelemetry-sdk and an exporter package — not a dependency of ldf itself):

import logging
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler
from opentelemetry.sdk._logs.export import BatchLogRecordProcessor
from opentelemetry.exporter.otlp.proto.http._log_exporter import OTLPLogExporter
from opentelemetry.sdk.resources import Resource

provider = LoggerProvider(resource=Resource.create({"service.name": "my-app"}))
provider.add_log_record_processor(BatchLogRecordProcessor(OTLPLogExporter()))
logging.getLogger().addHandler(LoggingHandler(logger_provider=provider))

Or, for local development, just logging.basicConfig(level=logging.INFO).

Token usage tracking

Token counts are accumulated on each Task's underlying LLM client (llm.input_token, llm.output_token) as calls complete. This works in both non-streaming and streaming modes — in streaming mode, token usage is captured from the final chunk via stream_options={"include_usage": True}. Task, MultiTask, and Workflow each expose a get_token_usage() method returning:

{
    "input_token_usage": ...,
    "output_token_usage": ...,
    "total_token_usage": ...,
}

For MultiTask and Workflow, this is a recursive sum over self.tasks — so a Workflow containing nested Workflows or MultiTasks reports combined usage for the whole tree:

await workflow.run()
print(workflow.get_token_usage())
# {"input_token_usage": 128, "output_token_usage": 342, "total_token_usage": 470}

Function doesn't implement get_token_usage, since it wraps a plain callable with no LLM attached. The aggregation in MultiTask/Workflow skips any task that lacks the method (via hasattr), so mixing in Functions or plain callables just contributes zero rather than raising an error.

Configuration

The LLM client is configured via environment variables or defined constructor arguments, so you can personalize it for each task:

Variable Default
LLM_MODEL gemma-26b-q5
LLM_BASE_URL http://localhost:8080/v1
LLM_API_KEY (empty)
LLM_SYSTEM_PROMPT You are a helpful assistant.
LLM_TEMPERATURE 0.7
LLM_MAX_TOKENS 1000
LLM_TIMEOUT 60 (seconds)

Constructor arguments take precedence over environment variables:

from ldf import LLM, Task

llm = LLM(model="gpt-4o-mini", temperature=0.2)
task = Task("greeting", "Hello!", llm=llm)

Sharing an HTTP connection pool across tasks

By default, each LLM instance builds its own openai.AsyncOpenAI client — fine for a script, but wasteful under load (e.g. many requests/tasks per process each opening their own connection pool). Pass a pre-built client to reuse one connection pool across many independently-configured LLM/Task instances:

import openai
from ldf import LLM, Task

shared_client = openai.AsyncOpenAI(api_key="...", base_url="...", timeout=60)

fast_task = Task("route", "...", llm=LLM(client=shared_client, model="gpt-4o-mini", temperature=0.0))
smart_task = Task("answer", "...", llm=LLM(client=shared_client, model="gpt-4o", system_prompt="Be thorough."))

Both tasks share the same underlying HTTP connections, but keep fully independent config and independent token counters — get_token_usage() is tracked per LLM instance, not per client, so sharing the connection pool never mixes up one task's usage with another's. If you don't pass client, you get today's behavior (a dedicated client per LLM instance) — it's opt-in, not a requirement.

Tip: share the client, not the LLM instance itself. Passing the same LLM object to multiple Tasks is unsafe — Task.__init__ sets llm.system_prompt directly, so two tasks sharing one LLM with different system_prompts will clobber each other's setting. Give each Task its own LLM(client=shared_client, ...) wrapper, as shown above, and only the connection pool gets shared.

MCP (Model Context Protocol) Integration

MCP parses MCP server configurations and creates LDF Tool instances for each tool discovered on the servers. It uses only stdlib (json, subprocess, select, urllib) — no external dependencies.

Input formats

MCP accepts three input formats:

  • File path to an mcp.json config file
  • JSON string with server configuration
  • URL to an MCP server (HTTP transport)

Transports

Two transports are supported:

  • stdio: spawns a local process, communicates via newline-delimited JSON-RPC over stdin/stdout
  • HTTP: sends JSON-RPC POST requests to a remote URL, tracks Mcp-Session-Id headers

Config format

The config follows the standard mcpServers schema used by Claude Code, Cursor, and other tools:

{
  "mcpServers": {
    "filesystem": {
      "command": "npx",
      "args": ["-y", "@modelcontextprotocol/server-filesystem", "/tmp"],
      "env": {}
    },
    "sentry": {
      "url": "https://mcp.sentry.dev/sse",
      "headers": {
        "Authorization": "Bearer ${SENTRY_TOKEN}"
      }
    }
  }
}

Also supports the servers key (used by VS Code) as an alternative to mcpServers.

Usage

from ldf import MCP, Agent

mcp = MCP("path/to/mcp.json")
tools = mcp.parse()

for tool in tools:
    print(f"  {tool.name}: {tool.description[:80]}")

agent = Agent(
    name="my_agent",
    user_input="Read /tmp/example.txt and search for errors",
    tools=tools,
    max_attempts=5,
    streaming=True,
)

result = await agent.run()

Or pass the MCP instance directly to Agent:

mcp = MCP("path/to/mcp.json")
agent = Agent(name="my_agent", user_input="...", tools=mcp)
result = await agent.run()

When using MCP with Agent, connections are closed automatically after run(). When using MCP standalone, close() is called automatically on garbage collection (__del__), but you can still call it explicitly if needed.

parse() performs the MCP handshake (initializenotifications/initializedtools/list) for each server and creates an LDF Tool for each discovered tool. Each tool's function callable handles parameter parsing and tools/call invocation via the appropriate transport.

close() terminates any spawned stdio processes. It's called automatically when used with Agent (after run()) or on garbage collection (__del__). You only need to call it manually if using MCP standalone and want deterministic cleanup.

Architecture

ldf/
├── __init__.py
├── agent.py            # Agent: LLM routing + multi-tool orchestration loop
├── tool.py             # Tool: LLM-assisted function execution with parameter extraction + result analysis
├── mcp.py              # MCP: Model Context Protocol parser, creates Tools from MCP servers
└── core/
    ├── llm.py          # LLM: OpenAI-compatible async client with token tracking
    ├── task.py         # Task: single LLM invocation with iterations
    ├── multitask.py    # MultiTask: concurrent (parallel) task groups
    ├── workflow.py     # Workflow: named sequence of tasks/task groups (can nest workflows)
    └── function.py     # Function: wraps a plain callable with condition/context support

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

ldf_framework-0.1.0.tar.gz (45.0 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

ldf_framework-0.1.0-py3-none-any.whl (19.1 kB view details)

Uploaded Python 3

File details

Details for the file ldf_framework-0.1.0.tar.gz.

File metadata

  • Download URL: ldf_framework-0.1.0.tar.gz
  • Upload date:
  • Size: 45.0 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.11.15

File hashes

Hashes for ldf_framework-0.1.0.tar.gz
Algorithm Hash digest
SHA256 14814149a599b863acbb91534243c5666f0008e334dd2c76a1ac8bf57970bfb3
MD5 bb21908b62acd42a3f530ed56dbb88f5
BLAKE2b-256 5bf6d5d025184243f5e80b96f7d52f89c42cc7b15ff315ef8961798e94d199db

See more details on using hashes here.

File details

Details for the file ldf_framework-0.1.0-py3-none-any.whl.

File metadata

  • Download URL: ldf_framework-0.1.0-py3-none-any.whl
  • Upload date:
  • Size: 19.1 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.11.15

File hashes

Hashes for ldf_framework-0.1.0-py3-none-any.whl
Algorithm Hash digest
SHA256 0a9b57c7a6205ebe6c59f9c9f707253122cdb5d0c4fde7e01d249230977e1652
MD5 9305037fdff976f85bb57c052540e1b7
BLAKE2b-256 d9aa59a0ee90c1c2e263e90261b8d716c5d3b02078a08a5dbb891eae88b7a238

See more details on using hashes here.

Release history Release notifications | RSS feed

This release

0.1.0 This release

2 files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page