|
|
@@ -1,10 +1,7 @@
|
|
|
-import { BusEvent } from "@/bus/bus-event"
|
|
|
-import { Bus } from "@/bus"
|
|
|
import { Log } from "../util/log"
|
|
|
import { describeRoute, generateSpecs, validator, resolver, openAPIRouteHandler } from "hono-openapi"
|
|
|
import { Hono } from "hono"
|
|
|
import { cors } from "hono/cors"
|
|
|
-import { streamSSE } from "hono/streaming"
|
|
|
import { proxy } from "hono/proxy"
|
|
|
import { basicAuth } from "hono/basic-auth"
|
|
|
import z from "zod"
|
|
|
@@ -34,6 +31,7 @@ import { FileRoutes } from "./routes/file"
|
|
|
import { ConfigRoutes } from "./routes/config"
|
|
|
import { ExperimentalRoutes } from "./routes/experimental"
|
|
|
import { ProviderRoutes } from "./routes/provider"
|
|
|
+import { EventRoutes } from "./routes/event"
|
|
|
import { InstanceBootstrap } from "../project/bootstrap"
|
|
|
import { NotFoundError } from "../storage/db"
|
|
|
import type { ContentfulStatusCode } from "hono/utils/http-status"
|
|
|
@@ -251,6 +249,7 @@ export namespace Server {
|
|
|
.route("/question", QuestionRoutes())
|
|
|
.route("/provider", ProviderRoutes())
|
|
|
.route("/", FileRoutes())
|
|
|
+ .route("/", EventRoutes())
|
|
|
.route("/mcp", McpRoutes())
|
|
|
.route("/tui", TuiRoutes())
|
|
|
.post(
|
|
|
@@ -498,64 +497,6 @@ export namespace Server {
|
|
|
return c.json(await Format.status())
|
|
|
},
|
|
|
)
|
|
|
- .get(
|
|
|
- "/event",
|
|
|
- describeRoute({
|
|
|
- summary: "Subscribe to events",
|
|
|
- description: "Get events",
|
|
|
- operationId: "event.subscribe",
|
|
|
- responses: {
|
|
|
- 200: {
|
|
|
- description: "Event stream",
|
|
|
- content: {
|
|
|
- "text/event-stream": {
|
|
|
- schema: resolver(BusEvent.payloads()),
|
|
|
- },
|
|
|
- },
|
|
|
- },
|
|
|
- },
|
|
|
- }),
|
|
|
- async (c) => {
|
|
|
- log.info("event connected")
|
|
|
- c.header("X-Accel-Buffering", "no")
|
|
|
- c.header("X-Content-Type-Options", "nosniff")
|
|
|
- return streamSSE(c, async (stream) => {
|
|
|
- stream.writeSSE({
|
|
|
- data: JSON.stringify({
|
|
|
- type: "server.connected",
|
|
|
- properties: {},
|
|
|
- }),
|
|
|
- })
|
|
|
- const unsub = Bus.subscribeAll(async (event) => {
|
|
|
- await stream.writeSSE({
|
|
|
- data: JSON.stringify(event),
|
|
|
- })
|
|
|
- if (event.type === Bus.InstanceDisposed.type) {
|
|
|
- stream.close()
|
|
|
- }
|
|
|
- })
|
|
|
-
|
|
|
- // Send heartbeat every 10s to prevent stalled proxy streams.
|
|
|
- const heartbeat = setInterval(() => {
|
|
|
- stream.writeSSE({
|
|
|
- data: JSON.stringify({
|
|
|
- type: "server.heartbeat",
|
|
|
- properties: {},
|
|
|
- }),
|
|
|
- })
|
|
|
- }, 10_000)
|
|
|
-
|
|
|
- await new Promise<void>((resolve) => {
|
|
|
- stream.onAbort(() => {
|
|
|
- clearInterval(heartbeat)
|
|
|
- unsub()
|
|
|
- resolve()
|
|
|
- log.info("event disconnected")
|
|
|
- })
|
|
|
- })
|
|
|
- })
|
|
|
- },
|
|
|
- )
|
|
|
.all("/*", async (c) => {
|
|
|
const path = c.req.path
|
|
|
|