Hacktoberfest 2026:メンテナが10月に向けて印を付けた、オープンで初心者向けの issue。 Hacktoberfest の issue を見る

Sometimes when downloading files via streaming, the data does not reach the onData event.

オープン
#3,033 コメント 2 件 リアクション 0 件 担当者 0 名 GitHub で見る

メンテナーはふだん 2 日以内に返信

まだ誰も着手していません。

評価

難易度
4/5
見積もり時間
3〜5日
初心者へのやさしさ
30/100
issue の種類
バグ
明瞭さ
説明が足りない
活発さ
停滞
技術スタック
electron, grpc, node.js, typescript
領域
api, backend

調査の方向性

レポートに示されているエントリーポイント client.DownloadFile/client.DownloadLog と、call.on("data") および call.on("end") のハンドラーから始めます。受信した FileBytes が handleBackpressure、PassThrough、pipeline を通じてどのように処理されるかを追跡し、どこでチャンクが失われる可能性があるかを特定します。不完全なファイルを再現し、関連するストリームの動作をドキュメント化または修正して、回帰テストを追加できれば完了です。

索引モデルが issue の本文から書いたものです。

説明

package: @grpc/grpc-js
Problem description

Sometimes the data sent by the server does not enter the data event.

Reproduction steps

This is a streaming file download function. The server sends data to the client in chunks until IsLastBlock is encountered, indicating that data transmission is complete. However, sometimes on('data') is not triggered, causing some chunks to be lost, resulting in an incomplete file.

Environment
Additional context

code segment:

import log from "electron-log";
import fs from "fs";
import fsPromises from "fs/promises";
import path from "path";
import { PassThrough, pipeline } from "stream";
import { promisify } from "util";
import { createGrpcClient } from "../grpc/client.js";
import { DownloadChildMessage } from "../types/grpcChild.js";

const pipe = promisify(pipeline);

log.info("[Worker] Log Downloader Worker Started");

process.on("message", async (msg: DownloadChildMessage & { userDataPath: string; protoPath: string }) => {
    const { type, ip, port, payload, taskId, instance, userDataPath, protoPath } = msg;

    log.info(`[Worker] gRPC target: ${ip}:${port}`);

    // 添加下载超时设置
    const DOWNLOAD_TIMEOUT = 5 * 60 * 1000; // 5分钟超时
    let timeoutId = setTimeout(() => {
        log.error(`[Worker] Download timeout after ${DOWNLOAD_TIMEOUT}ms`);
        process.send?.({ success: false, error: "Download timeout" });
        process.exit(1);
    }, DOWNLOAD_TIMEOUT);

    try {
        const client = createGrpcClient(ip, port, protoPath);
        const input = { InputFileJson: JSON.stringify(payload) };
        const tempDir = path.join(userDataPath, `downloads/logs/${taskId}`);
        await fsPromises.mkdir(tempDir, { recursive: true });

        const zipPath = path.join(
            tempDir,
            `${generateZipId(instance)}${type === "app" ? "_app" : ""}.zip`
        );

        const writeStream = fs.createWriteStream(zipPath);
        
        // 监听写入流的错误事件
        writeStream.on('error', (err) => {
            log.error(`[Worker] Write stream error: ${err.message}`);
            pass.destroy(err); // 将错误传播到pipeline
        });

        const pass = new PassThrough({ highWaterMark: 1024 * 1024 }); // 设置合适的缓冲区大小

        let receivedChunks = 0;
        let receivedBytes = 0;
        let lastBlockReceived = false;

        // pipeline 自动处理流和背压
        const pipelinePromise = pipe(pass, writeStream)
            .then(() => {
                clearTimeout(timeoutId);
                if (!lastBlockReceived) {
                    log.warn(`[Worker] Pipeline completed but IsLastBlock was never received. Received ${receivedChunks} chunks, ${receivedBytes} bytes.`);
                }
                log.info(`[Worker] Download complete -> ${zipPath}`);
                process.send?.({ success: true, path: zipPath });
            })
            .catch(err => {
                clearTimeout(timeoutId);
                log.error(`[Worker] Pipeline failed: ${err.message}`);
                process.send?.({ success: false, error: err.message });
            })
            .finally(() => {
                // 确保进程退出
                setTimeout(() => process.exit(0), 1000);
            });

        const call = type === "app" ? client.DownloadFile(input) : client.DownloadLog(input);

        // 改进的背压处理机制
        const handleBackpressure = (chunk: Buffer) => {
            if (!pass.write(chunk)) {
                // 缓冲区满,暂停gRPC流
                call.pause();
                log.debug(`[Worker] Backpressure: pausing gRPC stream. Buffer full.`);
                
                // 等待drain事件再恢复
                pass.once("drain", () => {
                    log.debug(`[Worker] Backpressure relieved: resuming gRPC stream.`);
                    call.resume();
                });
            }
        };

        call.on("data", (res: any) => {
            receivedChunks++;
            const chunkSize = res.FileBytes?.length || 0;
            receivedBytes += chunkSize;
            
            log.info(`[Worker] Received chunk ${receivedChunks}: ${chunkSize} bytes, IsLastBlock: ${res.IsLastBlock}`);
            
            if (res.IsLastBlock) {
                lastBlockReceived = true;
                log.info(`[Worker] Final block received. Total: ${receivedChunks} chunks, ${receivedBytes} bytes`);
            }
            
            if (res.FileBytes?.length) {
                handleBackpressure(res.FileBytes);
            } else {
                log.warn(`[Worker] Received chunk with empty FileBytes`);
            }
        });

        call.on("end", () => {
            log.info("[Worker] gRPC stream ended, closing PassThrough");
            if (!lastBlockReceived) {
                log.warn(`[Worker] Stream ended but IsLastBlock=true was not received. Received ${receivedChunks} chunks.`);
            }
            pass.end(); // 正常结束流
        });

        call.on("error", (err: any) => {
            clearTimeout(timeoutId);
            log.error(`[Worker] gRPC error: ${err.message}`);
            pass.destroy(err); // 将错误传播到pipeline
        });

        call.on("status", (status: any) => {
            log.info(`[Worker] gRPC status: ${JSON.stringify(status)}`);
        });

        // 等待pipeline完成
        await pipelinePromise;

    } catch (err: any) {
        clearTimeout(timeoutId);
        log.error(`[Worker] Fatal error: ${err.message}`);
        process.send?.({ success: false, error: err.message });
        setTimeout(() => process.exit(1), 1000);
    }
});

function generateZipId(instance: any) {
    const { HostName, ServerName, Version, ServiceInstanceName, Id } = instance;
    return `${HostName}_${ServerName}_${Version}_${ServiceInstanceName}_${Id}`;
}

import grpc from "@grpc/grpc-js";
import protoLoader from "@grpc/proto-loader";

export function createGrpcClient(ip: string, port: number, protoPath: string) {
    const packageDefinition = protoLoader.loadSync(protoPath, {
        keepCase: true,
        longs: String,
        enums: String,
        defaults: true,
        oneofs: true,
    });

    const proto = grpc.loadPackageDefinition(packageDefinition) as any;

    const Service = proto.ServerMangerGrpc;

    return new Service(
        `${ip}:${port}`,
        grpc.credentials.createInsecure()
    );
}


import { fork } from "child_process";
import { app } from "electron";
import log from "electron-log";
import path, { dirname } from "path";
import { fileURLToPath } from "url";
import { getProtoPath } from "../grpc/protoPath.js";
import { DownloadChildMessage } from "../types/grpcChild.js";

const __filename = fileURLToPath(import.meta.url);
const __dirname = dirname(__filename);

export const downloadLogInChild = (
    ip: string,
    port: number,
    payload: any,
    taskId: string,
    instance: any,
    type: "log" | "app" = "log",
    protoPath?: string
): Promise<string> => {
    return new Promise((resolve, reject) => {
        const child = fork(path.join(__dirname, "logDownloaderWorker.js"));
        let resolved = false; // 防止多次解析

        // 设置通信超时
        const timeout = setTimeout(() => {
            if (!resolved) {
                log.error(`Download child process timeout for ${instance.HostName}`);
                child.kill('SIGKILL');
                reject(new Error("Child process timeout"));
                resolved = true;
            }
        }, 10 * 60 * 1000); // 10分钟超时

        child.on("message", (msg: { success: boolean; path?: string; error?: string }) => {
            if (resolved) return;

            clearTimeout(timeout);
            resolved = true;

            if (msg.success) {
                log.info(`Download successful for ${instance.HostName}: ${msg.path}`);
                resolve(msg.path!);
            } else {
                log.error(`Download failed for ${instance.HostName}: ${msg.error}`);
                reject(new Error(msg.error || "Download failed"));
            }

            // 温和关闭子进程
            setTimeout(() => {
                if (child.connected) {
                    child.disconnect();
                }
                if (!child.killed) {
                    child.kill();
                }
            }, 1000);
        });

        child.on("error", (err) => {
            if (resolved) return;

            clearTimeout(timeout);
            resolved = true;
            log.error(`Child process error for ${instance.HostName}: ${err.message}`);
            reject(err);
        });

        child.on("exit", (code, signal) => {
            if (code !== 0 && !resolved) {
                clearTimeout(timeout);
                resolved = true;
                log.error(`Child process exited abnormally for ${instance.HostName}: code ${code}, signal ${signal}`);
                reject(new Error(`Child process exited with code ${code}`));
            }
        });

        // 传入 Electron 路径
        const userDataPath = app.getPath("userData");
        const message: DownloadChildMessage & { userDataPath: string; protoPath: string } = {
            ip, port, payload, taskId, instance, type,
            userDataPath,
            protoPath: protoPath || getProtoPath()
        };

        // 添加发送重试机制
        const sendMessage = (attempt = 0) => {
            if (attempt >= 3) {
                log.error(`Failed to send message to child process after ${attempt} attempts`);
                reject(new Error("Cannot communicate with child process"));
                return;
            }

            if (child.connected) {
                child.send(message, (err) => {
                    if (err) {
                        log.warn(`Failed to send message to child (attempt ${attempt + 1}): ${err.message}`);
                        setTimeout(() => sendMessage(attempt + 1), 1000);
                    }
                });
            } else {
                setTimeout(() => sendMessage(attempt + 1), 1000);
            }
        };

        sendMessage();
    });
};

log info:

Image Image

Thank you very much, could you please help me figure out where the problem is?

主要言語
TypeScript
スター
4.8k
フォーク
717
平均マージ
1日 18時間
マージ済み PR(30日)
17

環境構築

はじめの一歩

  1. issue を最後まで読み、次にプロジェクトのコントリビューションガイドを読みます。
  2. 着手することを issue にコメントします — 二人が同じ作業をするのを防げます。
  3. リポジトリをフォークし、ブランチを切って変更します。
  4. issue 番号を参照したプルリクエストを送ります。

grpc/grpc-node のほかの issue

grpc/grpc-node の issue をすべて見る

似ている issue

TypeScript の issue をもっと見る

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。