feat: buffer in memory. bun http2 doesn't work

This commit is contained in:
zk
2026-05-19 00:47:35 +01:00
parent 28db23f995
commit 88c271dc68
5 changed files with 84 additions and 80 deletions
@@ -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<Myt
const client = createClient(ChunithmUser, CreateMytTransport());
const request = create(GetPlaylogRequestSchema, { profileApiId });
try {
for await (const item of client.getPlaylog(request)) {
yield item;
}
} catch (err) {
if (err instanceof ConnectError) {
log.error(
{ err, code: err.code },
`MYT gRPC error streaming Chunithm playlog for userID ${userID}`,
);
} else {
log.error(
{ err },
`Unexpected MYT error streaming Chunithm playlog for userID ${userID}`,
);
}
throw new ScoreImportFatalError(500, `Failed to get scores from MYT.`);
}
yield* await drainMytPlaylogStream(client.getPlaylog(request), log, {
gameLabel: "Chunithm",
userID,
});
}
export default function ParseMytChunithm(
@@ -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 { GetPlaylogRequestSchema, MaimaiUser } from "#proto/generated/maimai/user_pb";
import { create } from "@bufbuild/protobuf";
import { ConnectError, createClient } from "@connectrpc/connect";
import { createClient } from "@connectrpc/connect";
import type { MytMaimaiDxScore } from "./types";
@@ -19,25 +19,10 @@ async function* streamPlaylog(userID: integer, log: KtLogger): AsyncIterable<Myt
const client = createClient(MaimaiUser, CreateMytTransport());
const request = create(GetPlaylogRequestSchema, { profileApiId });
try {
for await (const item of client.getPlaylog(request)) {
yield item;
}
} catch (err) {
if (err instanceof ConnectError) {
log.error(
{ err, code: err.code },
`MYT gRPC error streaming maimai DX playlog for userID ${userID}`,
);
} else {
log.error(
{ err },
`Unexpected MYT error streaming maimai DX playlog for userID ${userID}`,
);
}
throw new ScoreImportFatalError(500, `Failed to get scores from MYT.`);
}
yield* await drainMytPlaylogStream(client.getPlaylog(request), log, {
gameLabel: "maimai DX",
userID,
});
}
export default async function ParseMytMaimaiDx(
@@ -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 { GetPlaylogRequestSchema, OngekiUser } from "#proto/generated/ongeki/user_pb";
import { create } from "@bufbuild/protobuf";
import { ConnectError, createClient } from "@connectrpc/connect";
import { createClient } from "@connectrpc/connect";
import type { MytOngekiScore } from "./types";
@@ -19,25 +19,10 @@ async function* streamPlaylog(userID: integer, log: KtLogger): AsyncIterable<Myt
const client = createClient(OngekiUser, CreateMytTransport());
const request = create(GetPlaylogRequestSchema, { profileApiId });
try {
for await (const item of client.getPlaylog(request)) {
yield item;
}
} catch (err) {
if (err instanceof ConnectError) {
log.error(
{ err, code: err.code },
`MYT gRPC error streaming Ongeki playlog for userID ${userID}`,
);
} else {
log.error(
{ err },
`Unexpected MYT error streaming Ongeki playlog for userID ${userID}`,
);
}
throw new ScoreImportFatalError(500, `Failed to get scores from MYT.`);
}
yield* await drainMytPlaylogStream(client.getPlaylog(request), log, {
gameLabel: "Ongeki",
userID,
});
}
export default async function ParseMytOngeki(
@@ -4,39 +4,39 @@ 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 { PlaylogRequestSchema, WaccaUser } from "#proto/generated/wacca/user_pb";
import { create } from "@bufbuild/protobuf";
import { ConnectError, createClient } from "@connectrpc/connect";
import { createClient } from "@connectrpc/connect";
import type { MytWaccaScore } from "./types";
import CreateMytWACCAClassHandler from "./class-handler";
async function* streamPlaylog(apiId: string, log: KtLogger): AsyncIterable<MytWaccaScore> {
async function* streamPlaylog(
apiId: string,
log: KtLogger,
userID: integer,
): AsyncIterable<MytWaccaScore> {
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",
@@ -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<T>(
stream: AsyncIterable<T>,
log: KtLogger,
context: MytPlaylogStreamContext,
): Promise<T[]> {
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;
}