Skip to main content
Glama
yasinyaman
by yasinyaman

Fusion

DuckDB-powered in-memory analytics engine with LLM tool support.

Fusion connects to PostgreSQL/MySQL databases through the Warp REST API, lazily loads the tables a query needs into DuckDB, and exposes 10 analytics tools to LLMs over MCP, a REST API and a Python SDK.

Features

  • 10 LLM tools — list_sources, describe_table, query_data, search_data, aggregate_data, create_view, list_views, refresh_view, load_table, cache_stats

  • Three access layers — MCP server (stdio), REST API (FastAPI), Python SDK; tool definitions in OpenAI function-calling and MCP formats

  • Lazy loading — connecting a source fetches metadata only; tables are pulled on first use

  • Smart transfer — a table too big for memory is read one slice at a time: only the columns and rows a query touches, streamed as Arrow, reused by later queries, and refused with concrete advice when even that is too much

  • Query pushdown — single-source queries on unloaded tables run on the source database

  • Cross-source federation — JOIN PostgreSQL and MySQL tables in one DuckDB query

  • SQL guardrails — only read-only queries (SELECT, CTEs, UNION/INTERSECT/EXCEPT) reach DuckDB; file/network functions are denied and DuckDB's external access is latched off

  • Query cache, materialized views, backups — LRU cache with TTL, mv_* tables with scheduled refresh, timestamped backups with retention

  • Resilient Warp client — pooled connections with retry/backoff behind a circuit breaker, SSRF-guarded URLs

  • Hexagonal architecture — pure domain and application layers, adapters for DuckDB, sqlglot, Warp, FastAPI and MCP; every port has an in-memory fake for tests

Related MCP server: mcp-data-gateway

Architecture

                 inbound adapters                       outbound adapters
   ┌──────────┐ ┌──────────┐ ┌──────────┐      ┌─────────────┐ ┌─────────────┐
   │ REST API │ │ MCP srv  │ │ CLI/SDK  │      │ DuckDBStore │ │ WarpSource  │
   └────┬─────┘ └────┬─────┘ └────┬─────┘      └──────┬──────┘ └──────┬──────┘
        │            │            │                   │               │
        ▼            ▼            ▼            AnalyticsStore     DataSource
   ┌──────────────────────────────────────┐   ┌────────────────────────────────┐
   │ application  (FusionApp)             │   │ ports (Protocols)              │
   │  ToolService · QueryService          │◄──┤  AnalyticsStore · DataSource   │
   │  SourceService · MaterializedView    │   │  SqlValidator · SqlAnalyzer    │
   │  BackupService · FetchPlanner        │   │  QueryCache · Scheduler        │
   └───────────────────┬──────────────────┘   └────────────────────────────────┘
                       ▼
   ┌──────────────────────────────────────┐   sqlglot policy · memory cache
   │ domain  (pure Python)                │   threading scheduler · Warp HTTP
   │  TableRef · RowSet · QueryResult     │   (pool + circuit breaker)
   │  SchemaCatalog · guardrail text rules│
   └──────────────────────────────────────┘

Dependencies point inward: domain imports nothing, ports only domain, application only domain + ports. Adapters implement the ports, and fusion.bootstrap.build_app() wires them together. An architecture test enforces the rule.

Installation

Python 3.12+ and uv (or pip):

uv sync --all-extras          # development checkout
pip install "fusion[all]"     # everything
pip install "fusion[rest]"    # REST API (FastAPI + uvicorn)
pip install "fusion[mcp]"     # MCP server
pip install "fusion[pandas]"  # DataFrame conversions

The core package depends only on duckdb, pyarrow, requests and sqlglot.

Quick Start

Python SDK

from fusion import Settings, build_app

app = build_app(Settings(memory_limit="4GB"))          # or Settings.from_env()
app.sources.connect("mydb", {
    "type": "warp",
    "base_url": "http://localhost:8000",
    "database": "mydb",
})

app.tools.list_sources()                                  # what is available (metadata only)
app.tools.query_data("SELECT * FROM mydb.orders LIMIT 10")  # loads mydb.orders on first use
app.tools.aggregate_data("mydb.orders", "status", "amount", "SUM")
app.views.create("daily_revenue",
                 "SELECT status, SUM(amount) AS total FROM mydb.orders GROUP BY status",
                 refresh="hourly")

result = app.query.sql("SELECT COUNT(*) AS n FROM mydb.orders")  # QueryResult
result.to_records(); result.to_markdown(); result.to_json()

app.close()

build_app accepts replacement adapters (store=, cache=, scheduler=, source_factory=, validator=, analyzer=), which is how the tests plug in fakes.

MCP Server (Claude Desktop / Cursor)

fusion-mcp --warp-url http://localhost:8000 --database mydb
fusion-mcp --warp-url http://localhost:8000 --auto-discover
{
  "mcpServers": {
    "fusion": {
      "command": "fusion-mcp",
      "args": ["--warp-url", "http://localhost:8000", "--database", "mydb"]
    }
  }
}

REST API

fusion-rest --warp-url http://localhost:8000 --auto-discover --port 9000

Swagger UI at http://localhost:9000/docs.

Endpoint

Method

Description

/health, /readiness

GET

Liveness / readiness

/sources

GET

Connected sources and tables

/tables/{source.table}/schema

GET

Table schema

/query

POST

Analytical SQL ({"sql": ...})

/search

POST

Filter search on a table

/aggregate

POST

GROUP BY aggregation

/views

GET/POST

List or create materialized views

/views/{name}/refresh

POST

Refresh a materialized view

/tables/{source.table}/load

POST

Explicitly load a table

/cache/stats

GET

Cache statistics

/tools, /tools/{tool_name}

GET/POST

Tool definitions and generic dispatch

/backup/list, /backup/create, /backup/stats

GET/POST/GET

Backups

Configuration comes from FUSION_* / WARP_URL environment variables (see .env.example); in production an API key is mandatory and requests are rate-limited per key.

OpenAI Function Calling

from fusion import Settings, build_app, get_openai_tools

app = build_app(Settings())
app.sources.connect("mydb", {"type": "warp", "base_url": "http://localhost:8000"})

tools = get_openai_tools()                      # pass to the Chat Completions API
result = app.tools.execute("query_data", {"sql": "SELECT ..."})   # when the model calls a tool

Tools

Tool

Description

list_sources

Connected sources and tables with row counts, source-side estimates, load state and loaded slices

describe_table

Table schema (columns, types, row count)

query_data

Run analytical SQL on DuckDB (read-only, max 100 rows)

search_data

Filter search (exact match or LIKE with %)

aggregate_data

GROUP BY aggregation (SUM, AVG, COUNT, MIN, MAX)

create_view

Create a materialized view from a SELECT query

list_views

List materialized views with refresh schedule

refresh_view

Manually refresh a materialized view

load_table

Load a table, or just the slice where / columns describe

cache_stats

Query cache hit rate and entry count

Every tool returns a JSON-serializable dict; failures come back as {"error": "..."}.

Working with tables larger than memory

Fusion is an in-memory engine, so it will not load a table it cannot hold. Instead of guessing, it reads what a query actually needs:

# Only the matching rows and columns travel; the slice stays for later queries.
app.query.sql("SELECT id, total FROM shop.orders WHERE status = 'paid'")

# Or load a slice explicitly
app.tools.load_table("shop.orders", where="status = 'paid'", columns=["id", "total"])

A query with nothing to narrow a big table by is refused, with its estimated size and the ways forward (add a WHERE, select fewer columns, join it to a small table on an equality key, or raise the limit). A join between a huge table and a small one fetches only the rows whose key appears in the small one.

Slicing needs a source that can filter; Warp does. With Warp >= 0.10 the rows arrive as an Arrow IPC stream with their real types, and table sizes come from GET /schema without counting rows. Against Warp 0.9 everything still works through the paged list endpoint.

Setting

Default

Meaning

FUSION_FULL_LOAD_MAX_ROWS

500000

Largest table loaded whole when a query has no usable filter

FUSION_SLICE_MAX_ROWS

500000

Largest single slice

FUSION_SLICE_BUDGET_ROWS

2000000

Rows kept across all slices before the least used are evicted

FUSION_SEMI_JOIN_MAX_KEYS

50000

Most join keys passed to the source

FUSION_IN_CHUNK_SIZE

1000

Keys per request when passing them

FUSION_REFRESH_CONFIG

(none)

Per-table incremental refresh, as JSON

Warp Setup

git clone https://github.com/yasinyaman/warp.git
cd warp
docker compose up -d

Fusion authenticates with WARP_API_KEY, sent in the header Warp expects (X-API-Key by default, FUSION_WARP_API_KEY_HEADER to change it). Raw SQL pushdown is used only when Warp advertises it, so a Warp with enable_raw_query: false (the default, and mandatory in production) is fully supported through slices.

Incremental refresh keeps a table up to date without re-reading it:

export FUSION_REFRESH_CONFIG='{"shop.orders": {"watermark_column": "updated_at", "key_columns": ["id"]}}'

Project Structure

fusion/
├── __init__.py                 # public SDK: Settings, build_app, FusionApp, models, errors
├── bootstrap.py                # composition root (build_app, default_discovery)
├── domain/                     # pure Python: models, catalog, slices, query_shape, policy, identifiers, sql_text, views, errors
├── ports/                      # Protocols: DataSource, AnalyticsStore, SqlValidator/Analyzer, QueryCache, Scheduler
├── application/                # Settings, FetchPlanner, Source/Query/View/Backup/Tool services, SemiJoinExecutor, FusionApp
├── adapters/
│   ├── outbound/               # duckdb_store, sqlglot_policy, memory_cache, threading_scheduler, registry
│   │   └── warp/               # http (pool + circuit breaker + SSRF guard), capabilities, streams, source, discovery
│   └── inbound/
│       ├── rest/               # FastAPI app, routes, middleware (auth, logging), rate limit
│       ├── mcp/                # MCPServer adapter
│       ├── cli/                # fusion-rest, fusion-mcp entry points
│       └── sdk/                # pandas / Arrow conversions (optional)
└── observability/              # logging setup and formatters
tests/                          # domain, ports (contract tests), adapters, application, inbound, e2e
demo/demo.py                    # in-process demo over synthetic data

Development

uv sync --all-extras
uv run pytest                    # 900+ tests, coverage gate 80%
uv run ruff check fusion tests demo && uv run ruff format --check fusion tests demo
uv run mypy fusion               # strict on domain/ports/application
uv run python -m demo.demo       # demo with synthetic data (--scale 0.1 for a quick run)

Requirements

  • Python 3.12+

  • DuckDB 1.5+

  • Warp (data source gateway)

License

Apache 2.0 — see LICENSE.

Related MCP Connectors

Related MCP Servers

  • A
    license
    Not graded
    quality
    D
    maintenance
    Enables AI agents to query a PostgreSQL database through a small set of controlled, read-only tools for schema inspection, row lookup, and aggregate statistics.
    1
    MIT
  • A
    license
    Not graded
    quality
    D
    maintenance
    Enables AI agents to interact with PostgreSQL or MySQL databases using natural language. Supports SQL queries, schema discovery, and pre-built aggregations without writing SQL.
    233 npm
    MIT