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
"""
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
/**
* 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
"""
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
/**
* 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);