feat(tui): log event stream connections (#35973)
This commit is contained in:
parent
855a9569b5
commit
b642f9c7bc
3 changed files with 56 additions and 5 deletions
|
|
@ -7,9 +7,19 @@ import { createSimpleContext } from "./helper"
|
||||||
import { useLog } from "./log"
|
import { useLog } from "./log"
|
||||||
|
|
||||||
export type SDKConnectionStatus = "connected" | "connecting" | "reconnecting"
|
export type SDKConnectionStatus = "connected" | "connecting" | "reconnecting"
|
||||||
|
export type SDKConnectionEvent = {
|
||||||
|
readonly type: "client.connection"
|
||||||
|
readonly created: number
|
||||||
|
readonly data: {
|
||||||
|
readonly status: "connecting" | "connected" | "disconnected" | "reconnecting"
|
||||||
|
readonly attempt: number
|
||||||
|
readonly error?: string
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
type SDKEventMap = { [Type in V2Event["type"]]: Extract<V2Event, { type: Type }> }
|
type SDKEventMap = { [Type in V2Event["type"]]: Extract<V2Event, { type: Type }> }
|
||||||
const connectTimeout = 2_000
|
const connectTimeout = 2_000
|
||||||
|
const connectionHistoryLimit = 50
|
||||||
|
|
||||||
export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
||||||
name: "SDK",
|
name: "SDK",
|
||||||
|
|
@ -20,8 +30,9 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
||||||
// Stops and starts the managed service; present only in service mode.
|
// Stops and starts the managed service; present only in service mode.
|
||||||
reload?: () => Promise<void>
|
reload?: () => Promise<void>
|
||||||
}) => {
|
}) => {
|
||||||
const log = useLog()
|
const log = useLog({ component: "sdk" })
|
||||||
const abort = new AbortController()
|
const abort = new AbortController()
|
||||||
|
const history: SDKConnectionEvent[] = []
|
||||||
let client = props.client
|
let client = props.client
|
||||||
let api = props.api
|
let api = props.api
|
||||||
const events = createGlobalEmitter<SDKEventMap>()
|
const events = createGlobalEmitter<SDKEventMap>()
|
||||||
|
|
@ -35,6 +46,11 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
||||||
})
|
})
|
||||||
let stream: AbortController | undefined
|
let stream: AbortController | undefined
|
||||||
|
|
||||||
|
function record(status: SDKConnectionEvent["data"]["status"], attempt: number, error?: string) {
|
||||||
|
history.push({ type: "client.connection", created: Date.now(), data: { status, attempt, error } })
|
||||||
|
if (history.length > connectionHistoryLimit) history.shift()
|
||||||
|
}
|
||||||
|
|
||||||
function start() {
|
function start() {
|
||||||
stream?.abort()
|
stream?.abort()
|
||||||
const controller = new AbortController()
|
const controller = new AbortController()
|
||||||
|
|
@ -54,6 +70,8 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
||||||
)
|
)
|
||||||
controller.signal.addEventListener("abort", cancel, { once: true })
|
controller.signal.addEventListener("abort", cancel, { once: true })
|
||||||
const error = await (async () => {
|
const error = await (async () => {
|
||||||
|
record(attempt === 0 ? "connecting" : "reconnecting", attempt)
|
||||||
|
log.info("event stream connecting", { attempt })
|
||||||
const response = await client.v2.event.subscribe({
|
const response = await client.v2.event.subscribe({
|
||||||
signal: connection.signal,
|
signal: connection.signal,
|
||||||
sseMaxRetryAttempts: 0,
|
sseMaxRetryAttempts: 0,
|
||||||
|
|
@ -69,7 +87,9 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
||||||
if (first.value.type !== "server.connected")
|
if (first.value.type !== "server.connected")
|
||||||
return new Error("Event stream did not start with server.connected")
|
return new Error("Event stream did not start with server.connected")
|
||||||
clearTimeout(timeout)
|
clearTimeout(timeout)
|
||||||
|
record("connected", attempt)
|
||||||
attempt = 0
|
attempt = 0
|
||||||
|
log.info("event stream connected")
|
||||||
events.emit(first.value.type, first.value)
|
events.emit(first.value.type, first.value)
|
||||||
setConnection({ status: "connected", attempt: 0, error: undefined })
|
setConnection({ status: "connected", attempt: 0, error: undefined })
|
||||||
connected()
|
connected()
|
||||||
|
|
@ -93,6 +113,12 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
||||||
})
|
})
|
||||||
if (abort.signal.aborted || controller.signal.aborted) return
|
if (abort.signal.aborted || controller.signal.aborted) return
|
||||||
attempt += 1
|
attempt += 1
|
||||||
|
const message = error instanceof Error ? error.message : String(error)
|
||||||
|
record("disconnected", attempt, message)
|
||||||
|
log.info("event stream disconnected", {
|
||||||
|
attempt,
|
||||||
|
error: message,
|
||||||
|
})
|
||||||
// Re-resolve the transport before retrying: the server may have
|
// Re-resolve the transport before retrying: the server may have
|
||||||
// moved (service restarted on a new port) or need starting. Static
|
// moved (service restarted on a new port) or need starting. Static
|
||||||
// transports (--server, standalone) resolve to the same address.
|
// transports (--server, standalone) resolve to the same address.
|
||||||
|
|
@ -107,7 +133,7 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
||||||
setConnection({
|
setConnection({
|
||||||
status: "reconnecting",
|
status: "reconnecting",
|
||||||
attempt,
|
attempt,
|
||||||
error: error instanceof Error ? error.message : String(error),
|
error: message,
|
||||||
})
|
})
|
||||||
await wait(250, controller.signal)
|
await wait(250, controller.signal)
|
||||||
}
|
}
|
||||||
|
|
@ -143,6 +169,11 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
||||||
error() {
|
error() {
|
||||||
return connection.error
|
return connection.error
|
||||||
},
|
},
|
||||||
|
internal: {
|
||||||
|
history() {
|
||||||
|
return history.slice()
|
||||||
|
},
|
||||||
|
},
|
||||||
},
|
},
|
||||||
reload: props.reload,
|
reload: props.reload,
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -781,10 +781,19 @@ export function Session() {
|
||||||
)
|
)
|
||||||
: await (async () => {
|
: await (async () => {
|
||||||
if (options.debug) {
|
if (options.debug) {
|
||||||
const events: unknown[] = []
|
const events: { readonly created: number }[] = []
|
||||||
for await (const event of sdk.api.session.log({ sessionID: sessionData.id, follow: false })) {
|
for await (const event of sdk.api.session.log({ sessionID: sessionData.id, follow: false })) {
|
||||||
if (event.type !== "log.synced") events.push(event)
|
if (event.type !== "log.synced") events.push(event)
|
||||||
}
|
}
|
||||||
|
// Durable events stay in aggregate order even when their wall-clock timestamps differ.
|
||||||
|
sdk.connection.internal.history().forEach((event) => {
|
||||||
|
const index = events.findIndex((item) => item.created > event.created)
|
||||||
|
if (index === -1) {
|
||||||
|
events.push(event)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
events.splice(index, 0, event)
|
||||||
|
})
|
||||||
return JSON.stringify({ info: sessionData, events }, null, 2) + EOL
|
return JSON.stringify({ info: sessionData, events }, null, 2) + EOL
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -108,7 +108,9 @@ function Probe(props: {
|
||||||
describe("useEvent", () => {
|
describe("useEvent", () => {
|
||||||
test("logs only durable events", async () => {
|
test("logs only durable events", async () => {
|
||||||
const logs: Array<{ message: string; tags: Readonly<Record<string, unknown>> }> = []
|
const logs: Array<{ message: string; tags: Readonly<Record<string, unknown>> }> = []
|
||||||
const { app, emit, seen } = await mount(undefined, (_level, message, tags) => logs.push({ message, tags }))
|
const { app, emit, seen } = await mount(undefined, (_level, message, tags) => {
|
||||||
|
if (message === "event") logs.push({ message, tags })
|
||||||
|
})
|
||||||
const durable = event(
|
const durable = event(
|
||||||
{
|
{
|
||||||
id: "evt_renamed",
|
id: "evt_renamed",
|
||||||
|
|
@ -128,7 +130,7 @@ describe("useEvent", () => {
|
||||||
expect(logs).toEqual([
|
expect(logs).toEqual([
|
||||||
{
|
{
|
||||||
message: "event",
|
message: "event",
|
||||||
tags: { type: "session.renamed", aggregateID: "ses_test", seq: 1 },
|
tags: { component: "sdk", type: "session.renamed", aggregateID: "ses_test", seq: 1 },
|
||||||
},
|
},
|
||||||
])
|
])
|
||||||
} finally {
|
} finally {
|
||||||
|
|
@ -202,6 +204,15 @@ describe("useEvent", () => {
|
||||||
|
|
||||||
expect(sdk.client).toBe(replacement.client)
|
expect(sdk.client).toBe(replacement.client)
|
||||||
expect(sdk.api).toBe(replacement.api)
|
expect(sdk.api).toBe(replacement.api)
|
||||||
|
const history = sdk.connection.internal.history()
|
||||||
|
expect(history.map((event) => [event.data.status, event.data.attempt])).toEqual([
|
||||||
|
["connecting", 0],
|
||||||
|
["connected", 0],
|
||||||
|
["disconnected", 1],
|
||||||
|
["reconnecting", 1],
|
||||||
|
["connected", 1],
|
||||||
|
])
|
||||||
|
expect(history.every((event) => Number.isFinite(event.created))).toBe(true)
|
||||||
} finally {
|
} finally {
|
||||||
app.renderer.destroy()
|
app.renderer.destroy()
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue