Home/Templates

Clone a template

Each starter is one language. Copy that folder, set GRAPHINGEST_API_KEY, and run the one command on the card.

Python

vercel-cron-to-graphingest

The timeout escape hatch. The cron route returns in under a second.

graphingest/vercel-cron-to-graphingest · python pipeline.py

Read the guide

"""
Nightly sync that is too long for a Vercel function.

Deploy this file, then point app/api/cron/sync/route.ts at the flow id.

    pip install -r requirements.txt
    python pipeline.py
"""

from graphingest import node, graph, deploy, RetryPolicy
import requests


@node(name="sync-source", cache_ttl=1800, max_retries=3)
def sync_source(source: dict) -> dict:
    api_url = source.get("api_url") or source.get("apiUrl")
    resp = requests.get(api_url, timeout=30)
    data = resp.json()
    return {
        "source": source["name"],
        "records_synced": len(data) if isinstance(data, list) else 1,
    }


@graph(
    name="nightly-sync",
    timeout_seconds=7200,
    retry_policy=RetryPolicy(max_retries=2, delay_seconds=30, backoff_factor=2),
)
def nightly_sync(sources: list[dict]):
    results = sync_source.map(sources)
    total = sum(r["records_synced"] for r in results)
    return {"sources_synced": len(results), "total_records": total}


if __name__ == "__main__":
    deploy()
    result = nightly_sync([
        {"name": "users", "api_url": "https://jsonplaceholder.typicode.com/users"},
        {"name": "posts", "api_url": "https://jsonplaceholder.typicode.com/posts"},
    ])
    print(result)

TypeScript

ts-vercel-cron

TypeScript nightly job for a Vercel schedule. The cron route returns in under a second.

graphingest/ts-vercel-cron · npm run deploy

Read the guide

/**
 * Nightly sync that is too long for a Vercel function.
 * TypeScript starter for https://www.graphingest.io
 *
 *   npm install
 *   npm run deploy
 *
 * Then point app/api/cron/sync/route.ts at the flow id
 * from https://www.graphingest.io/flows
 */

import { node, graph, deploy } from "graphingest";

type Source = { name: string; apiUrl?: string; api_url?: string };
type SyncResult = { source: string; recordsSynced: number };

const syncSource = node(
  { name: "sync-source", cacheTtl: 1800, maxRetries: 3 },
  async (source: Source): Promise<SyncResult> => {
    const apiUrl = source.apiUrl ?? source.api_url;
    if (!apiUrl) throw new Error("Each source needs apiUrl");
    const resp = await fetch(apiUrl);
    const data = await resp.json();
    return {
      source: source.name,
      recordsSynced: Array.isArray(data) ? data.length : 1,
    };
  }
);

const nightlySync = graph(
  {
    name: "nightly-sync",
    timeoutMs: 7_200_000,
    retryPolicy: { maxRetries: 2, delayMs: 30_000, backoffFactor: 2 },
  },
  async (sources: Source[]) => {
    const results = (await syncSource.map(sources)) as SyncResult[];
    const total = results.reduce((sum, row) => sum + row.recordsSynced, 0);
    return { sourcesSynced: results.length, totalRecords: total };
  }
);

await deploy();

const result = await nightlySync([
  { name: "users", apiUrl: "https://jsonplaceholder.typicode.com/users" },
  { name: "posts", apiUrl: "https://jsonplaceholder.typicode.com/posts" },
]);
console.log(result);

Python

langgraph-fanout

LangGraph research agents in parallel, five at a time, with a throttle on the next batch.

graphingest/langgraph-fanout · python langgraph_fanout.py

Read the guide

"""
LangGraph fan-out with a rate limit.

One batch of questions becomes one agent per question. Agents run in
parallel, five at a time, so a long list cannot fire every model call
at once. ThrottlePolicy and ConcurrencyPolicy gate the graph itself:
the next batch waits until the window and a free slot allow it.

Run:
    pip install -r requirements.txt
    python langgraph_fanout.py

Swap `reason` for a real ChatOpenAI call when you are ready to spend tokens.
The fan-out and the limits stay the same.
"""

from graphingest import graph, deploy, ThrottlePolicy, ConcurrencyPolicy
from graphingest.langgraph import agent_node, AgentConfig


def build_researcher(config: AgentConfig):
    from langgraph.graph import StateGraph, END
    from typing import TypedDict

    class State(TypedDict):
        messages: list

    def reason(state: State) -> dict:
        question = state["messages"][-1]["content"]
        # Replace this with ChatOpenAI(model=config.model).invoke(...)
        answer = f"[{config.model}] notes on: {question}"
        return {"messages": [{"role": "assistant", "content": answer}]}

    builder = StateGraph(State)
    builder.add_node("reason", reason)
    builder.set_entry_point("reason")
    builder.add_edge("reason", END)
    return builder.compile()


researcher = agent_node(
    name="researcher",
    graph_builder=build_researcher,
    config=AgentConfig(
        model="gpt-4o-mini",
        system_prompt="Answer in one paragraph.",
        max_iterations=6,
        stream_steps=False,
    ),
)


@graph(
    name="research-fanout",
    timeout_seconds=3600,
    throttle=ThrottlePolicy(limit=20, period_seconds=60),
    concurrency=ConcurrencyPolicy(limit=5, wait_timeout_seconds=180),
)
def research_fanout(queries: list[str]):
    width = 5
    results = []
    for start in range(0, len(queries), width):
        results.extend(researcher.map(queries[start : start + width]))
    return results


if __name__ == "__main__":
    deploy()
    answers = research_fanout([
        "What is durable execution?",
        "How do serverless function timeouts work?",
        "What is a sliding-window rate limit?",
    ])
    for answer in answers:
        print(answer)

TypeScript

ts-fanout

Research assistants in parallel, five at a time, with a throttle on the next batch.

graphingest/ts-fanout · npm run deploy

Read the guide

/**
 * TypeScript fan-out for https://www.graphingest.io
 *
 * Five assistants work at a time inside one batch. A new batch may
 * start at most 20 times a minute, and at most five batches run together.
 *
 *   npm install
 *   npm run deploy
 *
 * Watch the run at https://www.graphingest.io/runs
 * Replace the sentence in `researcher` with your model when you are ready.
 */

import { node, graph, deploy } from "graphingest";

const researcher = node(
  { name: "researcher", maxRetries: 2, timeoutSeconds: 120 },
  async (question: string) => {
    const answer = `[gpt-4o-mini] notes on: ${question}`;
    return { question, answer };
  }
);

const researchFanout = graph(
  {
    name: "research-fanout",
    timeoutMs: 3_600_000,
    throttle: { limit: 20, periodSeconds: 60 },
    concurrency: { limit: 5, waitTimeoutSeconds: 180 },
  },
  async (queries: string[]) => {
    const width = 5;
    const results: { question: string; answer: string }[] = [];
    for (let start = 0; start < queries.length; start += width) {
      const batch = (await researcher.map(queries.slice(start, start + width))) as {
        question: string;
        answer: string;
      }[];
      results.push(...batch);
    }
    return results;
  }
);

await deploy();

const answers = await researchFanout([
  "What is durable execution?",
  "How do serverless function timeouts work?",
  "What is a sliding-window rate limit?",
]);
for (const answer of answers) {
  console.log(answer);
}

Python

nightly-batch

One action across a long list. Continue from the row that failed.

graphingest/nightly-batch · python nightly_batch.py

"""
Process a list with .map(). Cached items are not redone.

    pip install -r requirements.txt
    python nightly_batch.py

A failed run resumes from the failed item:

    GRAPHINGEST_RUN_ID=<run-id> python resume.py
"""

from graphingest import node, graph, deploy, RetryPolicy


@node(name="process-item", cache_ttl=86400, max_retries=2, timeout_seconds=120)
def process_item(item: dict) -> dict:
    # Raise here to see a partial failure, then resume.py.
    return {"id": item["id"], "ok": True}


@graph(
    name="nightly-batch",
    timeout_seconds=7200,
    retry_policy=RetryPolicy(max_retries=2, delay_seconds=10, backoff_factor=2),
)
def nightly_batch(items: list[dict]):
    results = process_item.map(items)
    return {"processed": len(results), "ids": [r["id"] for r in results]}


if __name__ == "__main__":
    deploy()
    items = [{"id": f"item-{i}"} for i in range(12)]
    print(nightly_batch(items))

TypeScript

ts-nightly-batch

One action across a long list. Continue from the row that failed.

graphingest/ts-nightly-batch · npm run deploy

/**
 * Process a list with .map(). Cached items are not redone.
 * TypeScript starter for https://www.graphingest.io
 *
 *   npm install
 *   npm run deploy
 *
 * A failed run resumes from the failed item:
 *   set GRAPHINGEST_RUN_ID to the id from https://www.graphingest.io/runs
 *   npm run resume
 */

import { node, graph, deploy } from "graphingest";

type Item = { id: string };
type ItemResult = { id: string; ok: boolean };

const processItem = node(
  { name: "process-item", cacheTtl: 86400, maxRetries: 2, timeoutSeconds: 120 },
  async (item: Item): Promise<ItemResult> => {
    return { id: item.id, ok: true };
  }
);

const nightlyBatch = graph(
  {
    name: "nightly-batch",
    timeoutMs: 7_200_000,
    retryPolicy: { maxRetries: 2, delayMs: 10_000, backoffFactor: 2 },
  },
  async (items: Item[]) => {
    const results = (await processItem.map(items)) as ItemResult[];
    return { processed: results.length, ids: results.map((row) => row.id) };
  }
);

await deploy();

const items = Array.from({ length: 12 }, (_, i) => ({ id: `item-${i}` }));
console.log(await nightlyBatch(items));

TypeScript

ts-webhook-pipeline

A message arrives, the route returns a run id, and the slow work continues.

graphingest/ts-webhook-pipeline · npm run deploy

/**
 * TypeScript pipeline triggered by app/api/webhook/route.ts.
 *
 *   npm install
 *   npm run deploy
 */

import { node, graph, deploy } from "graphingest";

const handleEvent = node(
  { name: "handle-event", cacheTtl: 300, maxRetries: 3 },
  async (event: { id: string; type: string }) => {
    return { id: event.id, type: event.type, handled: true };
  }
);

const pipeline = graph(
  {
    name: "webhook-pipeline",
    timeoutMs: 3_600_000,
    retryPolicy: { maxRetries: 2, delayMs: 1_000, backoffFactor: 2 },
  },
  async (event: { id: string; type: string }) => {
    return handleEvent(event);
  }
);

await deploy();

const result = await pipeline({ id: "evt_demo", type: "invoice.paid" });
console.log(result);