Files
WeKnora/docreader/proto/docreader_pb2_grpc.py
T
wizardchen 7b1bb1054f feat(docreader): speed up scanned-PDF parsing, stream image results, isolate heavy async queues
Large scanned PDFs (hundreds of pages) were slow and fragile end-to-end.
This change addresses the parse, transport, and task-scheduling layers:

docreader (parse + transport):
- Parallelize per-page scanned rendering across processes (forkserver/fork),
  with serial fallback. ~4-7x faster on large scanned PDFs; pdfium is not
  thread-safe so we fan out across processes. Configurable via
  DOCREADER_PDF_RENDER_PARALLELISM.
- Add server-streaming ReadStream RPC: emit one meta frame then one frame per
  image, so documents with many page images are no longer capped by the unary
  gRPC message-size limit (a 874-page PDF produced ~193MiB of images, far over
  the 50MB cap) and memory is bounded on both ends. Unary Read is kept for
  backward compatibility; the Go production reader switches to ReadStream.

VLM:
- Make the VLM HTTP timeout configurable (VLM_HTTP_TIMEOUT_SECONDS) and raise
  the default 90s -> 180s so dense scanned-page OCR does not time out with
  "context deadline exceeded".

Async task queues:
- Isolate high-volume, model-heavy fan-out tasks into dedicated asynq queues so
  a single large document cannot saturate the shared worker pool and block
  user-facing document parsing:
    image:multimodal  -> "multimodal"
    chunk:extract     -> "graph"
    question:generation -> "question"
- Register the new queues in the server weight map and the cancel inspector's
  scanned-queue set (so cancelling a knowledge still purges its pending tasks).
2026-06-03 12:29:13 +08:00

190 lines
7.0 KiB
Python

# Generated by the gRPC Python protocol compiler plugin. DO NOT EDIT!
"""Client and server classes corresponding to protobuf-defined services."""
import grpc
import warnings
from docreader.proto import docreader_pb2 as docreader__pb2
GRPC_GENERATED_VERSION = '1.80.0'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
try:
from grpc._utilities import first_version_is_lower
_version_not_supported = first_version_is_lower(GRPC_VERSION, GRPC_GENERATED_VERSION)
except ImportError:
_version_not_supported = True
if _version_not_supported:
raise RuntimeError(
f'The grpc package installed is at version {GRPC_VERSION},'
+ ' but the generated code in docreader_pb2_grpc.py depends on'
+ f' grpcio>={GRPC_GENERATED_VERSION}.'
+ f' Please upgrade your grpc module to grpcio>={GRPC_GENERATED_VERSION}'
+ f' or downgrade your generated code using grpcio-tools<={GRPC_VERSION}.'
)
class DocReaderStub(object):
"""Missing associated documentation comment in .proto file."""
def __init__(self, channel):
"""Constructor.
Args:
channel: A grpc.Channel.
"""
self.Read = channel.unary_unary(
'/docreader.DocReader/Read',
request_serializer=docreader__pb2.ReadRequest.SerializeToString,
response_deserializer=docreader__pb2.ReadResponse.FromString,
_registered_method=True)
self.ReadStream = channel.unary_stream(
'/docreader.DocReader/ReadStream',
request_serializer=docreader__pb2.ReadRequest.SerializeToString,
response_deserializer=docreader__pb2.ReadStreamResponse.FromString,
_registered_method=True)
self.ListEngines = channel.unary_unary(
'/docreader.DocReader/ListEngines',
request_serializer=docreader__pb2.ListEnginesRequest.SerializeToString,
response_deserializer=docreader__pb2.ListEnginesResponse.FromString,
_registered_method=True)
class DocReaderServicer(object):
"""Missing associated documentation comment in .proto file."""
def Read(self, request, context):
"""Missing associated documentation comment in .proto file."""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def ReadStream(self, request, context):
"""ReadStream is the streaming counterpart of Read. It first emits one
ReadStreamResponse carrying the parse metadata (markdown / metadata /
error), then emits one message per image. This keeps every gRPC message
small so large scanned PDFs (hundreds of page images, far exceeding the
unary message-size cap) can be returned without RESOURCE_EXHAUSTED and
with bounded memory on both ends.
"""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def ListEngines(self, request, context):
"""Missing associated documentation comment in .proto file."""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def add_DocReaderServicer_to_server(servicer, server):
rpc_method_handlers = {
'Read': grpc.unary_unary_rpc_method_handler(
servicer.Read,
request_deserializer=docreader__pb2.ReadRequest.FromString,
response_serializer=docreader__pb2.ReadResponse.SerializeToString,
),
'ReadStream': grpc.unary_stream_rpc_method_handler(
servicer.ReadStream,
request_deserializer=docreader__pb2.ReadRequest.FromString,
response_serializer=docreader__pb2.ReadStreamResponse.SerializeToString,
),
'ListEngines': grpc.unary_unary_rpc_method_handler(
servicer.ListEngines,
request_deserializer=docreader__pb2.ListEnginesRequest.FromString,
response_serializer=docreader__pb2.ListEnginesResponse.SerializeToString,
),
}
generic_handler = grpc.method_handlers_generic_handler(
'docreader.DocReader', rpc_method_handlers)
server.add_generic_rpc_handlers((generic_handler,))
server.add_registered_method_handlers('docreader.DocReader', rpc_method_handlers)
# This class is part of an EXPERIMENTAL API.
class DocReader(object):
"""Missing associated documentation comment in .proto file."""
@staticmethod
def Read(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(
request,
target,
'/docreader.DocReader/Read',
docreader__pb2.ReadRequest.SerializeToString,
docreader__pb2.ReadResponse.FromString,
options,
channel_credentials,
insecure,
call_credentials,
compression,
wait_for_ready,
timeout,
metadata,
_registered_method=True)
@staticmethod
def ReadStream(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_stream(
request,
target,
'/docreader.DocReader/ReadStream',
docreader__pb2.ReadRequest.SerializeToString,
docreader__pb2.ReadStreamResponse.FromString,
options,
channel_credentials,
insecure,
call_credentials,
compression,
wait_for_ready,
timeout,
metadata,
_registered_method=True)
@staticmethod
def ListEngines(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(
request,
target,
'/docreader.DocReader/ListEngines',
docreader__pb2.ListEnginesRequest.SerializeToString,
docreader__pb2.ListEnginesResponse.FromString,
options,
channel_credentials,
insecure,
call_credentials,
compression,
wait_for_ready,
timeout,
metadata,
_registered_method=True)