mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-24 12:41:02 +02:00
Set end-of-stream cleanly
This commit is contained in:
parent
f1180ecda2
commit
e02b834b22
4 changed files with 14 additions and 18 deletions
|
|
@ -173,7 +173,14 @@ class LibraryResponseTranslator(MessageTranslator):
|
||||||
return result
|
return result
|
||||||
|
|
||||||
def from_response_with_completion(self, obj: LibrarianResponse) -> Tuple[Dict[str, Any], bool]:
|
def from_response_with_completion(self, obj: LibrarianResponse) -> Tuple[Dict[str, Any], bool]:
|
||||||
"""Returns (response_dict, is_final)"""
|
"""Returns (response_dict, is_final)
|
||||||
# For streaming responses, check end_of_stream to determine if this is the final message
|
|
||||||
is_final = getattr(obj, 'end_of_stream', True)
|
For chunked streaming responses (total_chunks > 0), completion is
|
||||||
|
determined by whether we've reached the final chunk.
|
||||||
|
For non-streaming responses (total_chunks = 0), always final.
|
||||||
|
"""
|
||||||
|
if obj.total_chunks > 0:
|
||||||
|
is_final = (obj.chunk_index >= obj.total_chunks - 1)
|
||||||
|
else:
|
||||||
|
is_final = True
|
||||||
return self.from_pulsar(obj), is_final
|
return self.from_pulsar(obj), is_final
|
||||||
|
|
|
||||||
|
|
@ -212,10 +212,6 @@ class LibrarianResponse:
|
||||||
# list-uploads response
|
# list-uploads response
|
||||||
upload_sessions: list[UploadSession] = field(default_factory=list)
|
upload_sessions: list[UploadSession] = field(default_factory=list)
|
||||||
|
|
||||||
# stream-document response - indicates final chunk in stream
|
|
||||||
# Default True so non-streaming operations are treated as final
|
|
||||||
end_of_stream: bool = True
|
|
||||||
|
|
||||||
# FIXME: Is this right? Using persistence on librarian so that
|
# FIXME: Is this right? Using persistence on librarian so that
|
||||||
# message chunking works
|
# message chunking works
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -656,9 +656,8 @@ class Librarian:
|
||||||
|
|
||||||
This is an async generator that yields document content in smaller chunks,
|
This is an async generator that yields document content in smaller chunks,
|
||||||
allowing memory-efficient processing of large documents. Each yielded
|
allowing memory-efficient processing of large documents. Each yielded
|
||||||
response includes chunk information and an end_of_stream flag.
|
response includes chunk_index and total_chunks for tracking progress.
|
||||||
|
Completion is determined by chunk_index reaching total_chunks - 1.
|
||||||
The final chunk will have end_of_stream=True.
|
|
||||||
"""
|
"""
|
||||||
logger.debug(f"Streaming document {request.document_id}")
|
logger.debug(f"Streaming document {request.document_id}")
|
||||||
|
|
||||||
|
|
@ -688,11 +687,8 @@ class Librarian:
|
||||||
# Fetch only the requested range
|
# Fetch only the requested range
|
||||||
chunk_content = await self.blob_store.get_range(object_id, offset, length)
|
chunk_content = await self.blob_store.get_range(object_id, offset, length)
|
||||||
|
|
||||||
is_last_chunk = (chunk_index == total_chunks - 1)
|
logger.debug(f"Streaming chunk {chunk_index + 1}/{total_chunks}, "
|
||||||
|
f"bytes {offset}-{offset + length} of {total_size}")
|
||||||
logger.debug(f"Streaming chunk {chunk_index}/{total_chunks}, "
|
|
||||||
f"bytes {offset}-{offset + length} of {total_size}, "
|
|
||||||
f"end_of_stream={is_last_chunk}")
|
|
||||||
|
|
||||||
yield LibrarianResponse(
|
yield LibrarianResponse(
|
||||||
error=None,
|
error=None,
|
||||||
|
|
@ -702,6 +698,5 @@ class Librarian:
|
||||||
total_chunks=total_chunks,
|
total_chunks=total_chunks,
|
||||||
bytes_received=offset + length,
|
bytes_received=offset + length,
|
||||||
total_bytes=total_size,
|
total_bytes=total_size,
|
||||||
end_of_stream=is_last_chunk,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -527,7 +527,6 @@ class Processor(AsyncProcessor):
|
||||||
type = "request-error",
|
type = "request-error",
|
||||||
message = str(e),
|
message = str(e),
|
||||||
),
|
),
|
||||||
end_of_stream = True,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
await self.librarian_response_producer.send(
|
await self.librarian_response_producer.send(
|
||||||
|
|
@ -541,7 +540,6 @@ class Processor(AsyncProcessor):
|
||||||
type = "unexpected-error",
|
type = "unexpected-error",
|
||||||
message = str(e),
|
message = str(e),
|
||||||
),
|
),
|
||||||
end_of_stream = True,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
await self.librarian_response_producer.send(
|
await self.librarian_response_producer.send(
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue