Loading...
Loading...
可长期运行的无服务器Node.js HTTP函数,部署到您的Neon分支上,自动注入DATABASE_URL,计算资源与数据就近运行。适用于用户想要托管API、带有长流式响应的AI Agent、WebSocket或服务器发送事件(SSE)服务器、Webhook处理程序、Discord机器人、MCP服务器,或任何可能在短时长的Lambda风格无服务器函数上超时的请求/响应工作负载——并且希望该函数能随数据库分支同步。触发场景包括“无服务器函数”、“部署API”、“长期运行函数”、“流式Agent”、“SSE服务器”、“WebSocket服务器”、“Webhook处理程序”、“MCP服务器”、“在我的数据库旁运行代码”、“不会超时的函数”、“函数日志”、“Neon Functions”和“Neon Compute”。
npx skill4agent add neondatabase/agent-skills neon-functionsneonneonnpx skills add neondatabase/agent-skills --skill neonus-east-2DATABASE_URLneon.tsneon devpgDATABASE_URLwaitUntilfetch(request)Responseexport default appDATABASE_URLneonneon.tsus-east-2us-east-2neon.tsneonneon.ts@neon/configpreview.functions// neon.ts
import { defineConfig } from "@neon/config/v1";
export default defineConfig({
preview: {
functions: {
todos: {
// slug格式:^[a-z0-9]{1,20}$ —— 小写字母/数字,无连字符
name: "todo api", // 仅为显示标签
source: "src/index.ts", // 相对于neon.ts的入口文件
},
},
},
});nameDATABASE_URL// src/index.ts
import { Hono } from "hono";
import { drizzle } from "drizzle-orm/node-postgres";
import { Pool } from "pg";
import { parseEnv } from "@neon/env";
import config from "../neon";
import { todos } from "./db/schema";
const env = parseEnv(config);
const pool = new Pool({ connectionString: env.postgres.databaseUrl, max: 5 });
const db = drizzle(pool);
const app = new Hono();
app.get("/", (c) => c.text("Neon + Hono + Drizzle"));
app.post("/todos", async (c) => {
const { text } = await c.req.json<{ text: string }>();
const [row] = await db.insert(todos).values({ text }).returning();
return c.json(row, 201);
});
app.get("/todos", async (c) => c.json(await db.select().from(todos)));
export default app;pgmaxparseEnv(config)parseEnvneon.tsconst { postgres } = parseEnv(config, ["DATABASE_URL"]); // 非非池化URL、认证信息等
const pool = new Pool({ connectionString: postgres.databaseUrl, max: 5 });neon dev # 热重载运行neon.ts中的所有函数;注入DATABASE_URL等环境变量
neon deploy # 使用esbuild打包、上传,并将neon.ts应用到关联分支neon.tsneon functions deploy <slug> --src src/index.ts--srcindex.tsindex.mjsindex.jsneon functions get <slug>invocation_urlhttps://<branch_id>-<slug>.compute.c-1.us-east-2.aws.neon.techneon functions list|get|deleteneon checkoutneon.tsneon deployneon.tspreview.functionsneon.tsneon.tssourcenameenvneonneon config status # 打印分支的实时配置(已部署的函数)
neon config plan # 试运行apply操作的差异
neon config apply # 打包并部署声明的函数(neon deploy是别名)neon.tsneon checkoutneon deployruntimebranchexport default defineConfig({
preview: {
functions: { todos: { name: "todo api", source: "src/index.ts" } },
},
branch: (branch) => ({
preview: { functions: { todos: { runtime: "nodejs24" } } },
}),
});| 变量名 | 说明 |
|---|---|
| 分支名称(如 |
| 池化连接字符串。适用于大多数查询。仅当分支包含Postgres时存在。 |
| 直接连接字符串。适用于迁移、 |
| 当分支启用Neon Auth时存在。 |
| 当分支启用Data API时存在。 |
AWS_*NEON_AI_GATEWAY_*neon-object-storageneon-ai-gatewayneon env pullneon-env runneon devNEON_BRANCHneon functions deploy--env KEY=VALUE--env KEY=neon.tsenvprocess.envfunctions: {
todos: {
name: "todo api",
source: "src/index.ts",
env: { RESEND_API_KEY: process.env.RESEND_API_KEY! },
},
}neon deploy --env .env.production.envneon env pulllinkcheckout--no-env-pullneon-env run -- <cmd>NEON_DATABASE_URLDATABASE_URLDATABASE_URL_UNPOOLEDLISTENNOTIFYpgpgimport { drizzle } from "drizzle-orm/node-postgres";
import { Pool } from "pg";
// 每个隔离实例创建一次;由该隔离实例处理的所有请求复用。
const pool = new Pool({ connectionString: process.env.DATABASE_URL, max: 5 });
const db = drizzle(pool);max5SIGINTSIGTERMwaitUntilwaitUntil@neon/functionswaitUntilSIGINTprocess.on("SIGINT", ...)^[a-z0-9]{1,20}$neon-ai-gatewaytoUIMessageStreamResponse浏览器 ──(Authorization: Bearer <JWT>)──▶ Neon Function (agent) ✅ 无宿主超时
浏览器 ──▶ 您的应用后端 ──▶ Neon Function ❌ 宿主切断流jwtnew DefaultChatTransport({ api: NEON_FUNCTION_URL, fetch })fetchAuthorization: Bearer <token>OPTIONSAccess-Control-Allow-Origin-Headers[!WARNING] Neon Function有一个公开HTTPS URL——任何人都可以访问。客户端→函数的直接调用意味着没有应用后端在前面进行访问控制,因此您必须自行对函数进行身份验证。在处理程序顶部验证JWT(如针对您应用的JWKS)、检查共享密钥/API密钥,或验证会话令牌,拒绝其他任何请求。切勿部署未认证的Agent。
// src/index.ts —— 在执行任何工作前验证调用者
import { createRemoteJWKSet, jwtVerify } from "jose";
const jwks = createRemoteJWKSet(new URL(`${process.env.AUTH_BASE_URL}/api/auth/jwks`));
export default {
async fetch(request: Request) {
if (request.method === "OPTIONS") return new Response(null, { status: 204, headers: cors(request) });
const auth = request.headers.get("authorization");
if (!auth?.toLowerCase().startsWith("bearer ")) {
return new Response("Unauthorized", { status: 401, headers: cors(request) });
}
try {
const { payload } = await jwtVerify(auth.slice(7), jwks, {
issuer: process.env.AUTH_BASE_URL,
audience: process.env.AUTH_BASE_URL,
});
const userId = payload.sub; // 将Agent限定为该用户
// ... 运行Agent,返回result.toUIMessageStreamResponse({ headers: cors(request) })
} catch {
return new Response("Unauthorized", { status: 401, headers: cors(request) });
}
},
};envimport { upgradeWebSocket } from "@neon/functions";
export default {
async fetch(req: Request): Promise<Response> {
if (req.headers.get("upgrade")?.toLowerCase() !== "websocket") {
return new Response("expected a websocket upgrade", { status: 426 });
}
const { socket, response } = upgradeWebSocket(req);
socket.addEventListener("message", (event) => socket.send(event.data));
return response;
},
};socketWebSocketaddEventListeneronopenonmessageoncloseonerrorCONNECTINGresponse101responseResponse101clone()new Response(res.body, res)Response401403404binaryType"arraybuffer""blob"event.datastringArrayBuffertypeof// src/index.ts
import { upgradeWebSocket } from "@neon/functions";
const clients = new Set<WebSocket>();
export default {
async fetch(request: Request): Promise<Response> {
if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") {
return new Response("WebSocket endpoint — connect with ?token=<jwt>");
}
const url = new URL(request.url);
const identity = await verifyToken(url.searchParams.get("token"));
if (!identity) return new Response("unauthorized", { status: 401 });
const { socket, response } = upgradeWebSocket(request);
clients.add(socket);
socket.addEventListener("close", () => clients.delete(socket));
socket.addEventListener("message", (event) => {
if (typeof event.data !== "string") return;
persist(identity.id, event.data); // 扩散到每个隔离实例——请查看下文
});
return response;
},
};{ protocol }Sec-WebSocket-Protocolsocket.protocolTypeErrorsocket.extensions""permessage-deflateupgradeWebSocketRequestc.req.raw// src/index.ts
import { Hono } from "hono";
import { upgradeWebSocket } from "@neon/functions";
const app = new Hono();
app.get("/", (c) => c.text("ok"));
app.get("/ws", async (c) => {
const identity = await verifyToken(c.req.query("token"));
if (!identity) return c.text("Unauthorized", 401);
const { socket, response } = upgradeWebSocket(c.req.raw);
socket.addEventListener("open", () => socket.send("welcome"));
socket.addEventListener("message", (event) => socket.send(`echo: ${event.data}`));
return response;
});
export default { fetch: (request: Request) => app.fetch(request) };WebSocketping()const HEARTBEAT_MS = 25_000; // 远低于代理空闲超时
const beat = setInterval(() => {
for (const socket of clients) {
if (socket.readyState === socket.OPEN) socket.send('{"type":"ping"}');
}
}, HEARTBEAT_MS);
beat.unref?.();clientsneon devpoolpgclientsSetlet lastId = 0;
const poller = setInterval(async () => {
if (clients.size === 0) return; // 此处无客户端 → 无查询 → 计算可缩容至零
const { rows } = await pool.query(
"SELECT id, payload FROM events WHERE id > $1 ORDER BY id",
[lastId],
);
for (const { id, payload } of rows) {
lastId = id;
for (const socket of clients) {
if (socket.readyState === socket.OPEN) socket.send(payload);
}
}
}, 1000);
poller.unref?.();serialbigserialLISTENNOTIFYLISTENNOTIFYimport { Pool, Client } from "pg";
const pool = new Pool({ connectionString: process.env.DATABASE_URL, max: 5 });
const CHANNEL = "chat_events";
// 每个隔离实例一个专用的直接连接,仅用于接收事件。
// 使用DATABASE_URL_UNPOOLED —— LISTEN需要真实会话,而非池化连接。
const listener = new Client({ connectionString: process.env.DATABASE_URL_UNPOOLED });
listener.connect().then(() => listener.query(`LISTEN ${CHANNEL}`));
listener.on("notification", (msg) => {
if (!msg.payload) return;
for (const socket of clients) {
if (socket.readyState === socket.OPEN) socket.send(msg.payload);
}
});
// 通过池化连接NOTIFY进行广播——每个隔离实例的监听器都会触发。
function broadcast(event: unknown) {
return pool.query("SELECT pg_notify($1, $2)", [CHANNEL, JSON.stringify(event)]);
}LISTENNOTIFYlet closed = false, retry = 0, timer: ReturnType<typeof setTimeout>;
async function connect() {
if (closed) return;
const token = await getToken(); // 每次尝试重新生成;短期有效
const ws = new WebSocket(`${WS_URL}?token=${encodeURIComponent(token)}`);
ws.onopen = () => { retry = 0; }; // 成功时重置退避
ws.onmessage = (e) => { /* 应用事件 */ };
ws.onclose = () => {
if (!closed) timer = setTimeout(connect, Math.min(1000 * 2 ** retry++, 15000));
};
ws.onerror = () => ws.close(); // 让onclose驱动重试
}
connect();fetchupgradeWebSocket?token=fetchResponseContent-Type: text/event-streamReadableStreamEventSource// src/index.ts —— 极简SSE端点
const encoder = new TextEncoder();
export default {
fetch: () =>
new Response(
new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(encoder.encode("data: hello\
\
"));
const t = setInterval(() => controller.enqueue(encoder.encode(": ping\
\
")), 25_000);
return () => clearInterval(t); // 客户端断开连接时触发
},
}),
{ headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache, no-transform" } },
),
};: ping\ \ SetenqueueEventSource?token=POSTGET/mcpfetch@modelcontextprotocol/sdk@hono/mcp/mcpconst transport = new StreamableHTTPTransport();
app.all("/mcp", async (c) => {
if (!mcpServer.isConnected()) await mcpServer.connect(transport);
return transport.handleRequest(c);
});mcporteradd-mcpneon logs query --branch production --source function --since 1h--branchneon--envneon.tsenv.mdAccept: text/markdown