From 88c271dc682ec2ce5075f9595f052c1ed1ba3a67 Mon Sep 17 00:00:00 2001 From: zk Date: Tue, 19 May 2026 00:47:35 +0100 Subject: [PATCH] feat: buffer in memory. bun http2 doesn't work --- .../import-types/api/myt-chunithm/parser.ts | 27 +++------- .../import-types/api/myt-maimaidx/parser.ts | 27 +++------- .../import-types/api/myt-ongeki/parser.ts | 27 +++------- .../import-types/api/myt-wacca/parser.ts | 34 ++++++------- .../common/api-myt/buffer-playlog-stream.ts | 49 +++++++++++++++++++ 5 files changed, 84 insertions(+), 80 deletions(-) create mode 100644 typescript/server/src/lib/score-import/import-types/common/api-myt/buffer-playlog-stream.ts diff --git a/typescript/server/src/lib/score-import/import-types/api/myt-chunithm/parser.ts b/typescript/server/src/lib/score-import/import-types/api/myt-chunithm/parser.ts index 2ed252851..1cb2be529 100644 --- a/typescript/server/src/lib/score-import/import-types/api/myt-chunithm/parser.ts +++ b/typescript/server/src/lib/score-import/import-types/api/myt-chunithm/parser.ts @@ -3,14 +3,14 @@ import type { ParserFunctionReturns } from "#lib/score-import/import-types/commo import type { EmptyObject } from "#utils/types"; import type { GamesForGroup, integer } from "tachi-common"; -import ScoreImportFatalError from "#lib/score-import/framework/score-importing/score-import-error"; +import { drainMytPlaylogStream } from "#lib/score-import/import-types/common/api-myt/buffer-playlog-stream"; import { CreateMytTransport, FetchMytTitleAPIID, } from "#lib/score-import/import-types/common/api-myt/traverse-api"; import { ChunithmUser, GetPlaylogRequestSchema } from "#proto/generated/chunithm/user_pb"; import { create } from "@bufbuild/protobuf"; -import { ConnectError, createClient } from "@connectrpc/connect"; +import { createClient } from "@connectrpc/connect"; import type { MytChunithmScore } from "./types"; @@ -19,25 +19,10 @@ async function* streamPlaylog(userID: integer, log: KtLogger): AsyncIterable { +async function* streamPlaylog( + apiId: string, + log: KtLogger, + userID: integer, +): AsyncIterable { const client = createClient(WaccaUser, CreateMytTransport()); const request = create(PlaylogRequestSchema, { apiId }); - try { - for await (const item of client.getPlaylog(request)) { - if (!item.info) { - log.warn(`Received WACCA playlog stream item with no info - skipping.`); - continue; - } + const items = await drainMytPlaylogStream(client.getPlaylog(request), log, { + gameLabel: "WACCA", + userID, + }); - yield item.info; - } - } catch (err) { - if (err instanceof ConnectError) { - log.error({ err, code: err.code }, `MYT gRPC error streaming WACCA playlog`); - } else { - log.error({ err }, `Unexpected MYT error streaming WACCA playlog`); + for (const item of items) { + if (!item.info) { + log.warn(`Received WACCA playlog stream item with no info - skipping.`); + continue; } - throw new ScoreImportFatalError(500, `Failed to get scores from MYT.`); + yield item.info; } } @@ -57,7 +57,7 @@ export default async function ParseMytWACCA( return { service: "MYT", - iterable: streamPlaylog(titleApiId, log), + iterable: streamPlaylog(titleApiId, log, userID), context: {}, classProvider, gameGroup: "wacca", diff --git a/typescript/server/src/lib/score-import/import-types/common/api-myt/buffer-playlog-stream.ts b/typescript/server/src/lib/score-import/import-types/common/api-myt/buffer-playlog-stream.ts new file mode 100644 index 000000000..3ce7acb47 --- /dev/null +++ b/typescript/server/src/lib/score-import/import-types/common/api-myt/buffer-playlog-stream.ts @@ -0,0 +1,49 @@ +import type { KtLogger } from "#lib/log/log"; + +import ScoreImportFatalError from "#lib/score-import/framework/score-importing/score-import-error"; +import { ConnectError } from "@connectrpc/connect"; +import type { integer } from "tachi-common"; + +export type MytPlaylogStreamContext = { + gameLabel: string; + userID?: integer; +}; + +/** + * Drain a MYT server-streaming GetPlaylog RPC before score import processes rows. + * + * Bun's HTTP/2 client closes the stream with "Premature close" if we yield each + * item and then do slow DB work before reading the next message. Draining first + * matches Node behaviour and the fast path of the myt-grpc-probe. + */ +export async function drainMytPlaylogStream( + stream: AsyncIterable, + log: KtLogger, + context: MytPlaylogStreamContext, +): Promise { + const items: T[] = []; + + try { + for await (const item of stream) { + items.push(item); + } + } catch (err) { + const userSuffix = context.userID !== undefined ? ` for userID ${context.userID}` : ""; + + if (err instanceof ConnectError) { + log.error( + { err, code: err.code }, + `MYT gRPC error streaming ${context.gameLabel} playlog${userSuffix}`, + ); + } else { + log.error({ err }, `Unexpected MYT error streaming ${context.gameLabel} playlog${userSuffix}`); + } + + throw new ScoreImportFatalError(500, `Failed to get scores from MYT.`); + } + + const userSuffix = context.userID !== undefined ? ` for userID ${context.userID}` : ""; + log.debug(`Buffered ${items.length} MYT ${context.gameLabel} playlog row(s)${userSuffix}`); + + return items; +}