feat: python sdk batch operations & smart batching - #3

Merged
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching
Mar 24, 2026
Merged

feat: python sdk batch operations & smart batching#3
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching

Conversation

@vieiralucas

@vieiralucasvieiralucas commented Mar 24, 2026

Copy link
Copy Markdown
Member

Summary

  • Add batch_enqueue() method for explicit multi-message RPCs on both sync and async clients
  • Add smart batching via BatchMode enum (AUTO/DISABLED) and Linger dataclass, routing enqueue() through a background batcher thread by default
  • Add delivery batching: consume stream transparently unpacks ConsumeResponse.messages repeated field with backward-compatible fallback to singular message field
  • Single-item optimization: when flushing 1 message, use singular Enqueue RPC to preserve error types (QueueNotFoundError vs BatchEnqueueError)
  • close() drains pending batched messages before disconnecting
  • Update proto to include BatchEnqueue RPC and ConsumeResponse.messages repeated field
  • 16 new unit tests (batcher mechanics, flush functions, batch mode types) + 9 integration tests (batch enqueue, smart batching modes)

Test plan

  • Unit tests for _flush_single and _flush_batch with mock stubs
  • Unit tests for AutoBatcher (single message, concurrent batching, close drains, stub update)
  • Unit tests for LingerBatcher (batch size trigger, linger timeout, close drains)
  • Unit tests for BatchMode, Linger, BatchEnqueueResult types
  • Integration tests for explicit batch_enqueue() (multiple messages, single message, consume verification)
  • Integration tests for async batch_enqueue()
  • Integration tests for smart batching modes (AUTO, DISABLED, LINGER, default-is-AUTO)
  • Existing integration tests still pass (enqueue/consume/ack lifecycle, nack redelivery, API key auth)
  • ruff lint clean, mypy clean (no new errors)

Generated with Claude Code


Summary by cubic

Add explicit batch_enqueue() to the Python SDK and enable smart batching by default (AUTO) for enqueue(). Also add batched delivery in consume() and update the proto with BatchEnqueue and ConsumeResponse.messages.

  • New Features

    • Client.batch_enqueue() and AsyncClient.batch_enqueue() send many messages in one RPC; return per-message BatchEnqueueResult.
    • Smart batching via BatchMode (AUTO default, DISABLED) and Linger; background batcher with single-item optimization; close() drains pending; reconnect updates batcher stub.
    • Delivery batching: consume() transparently iterates ConsumeResponse.messages with fallback to singular message.
    • Proto changes: add BatchEnqueue RPC and ConsumeResponse.messages; new BatchEnqueueError; public exports updated.
  • Migration

    • enqueue() now routes through a background batcher by default; disable with batch_mode=BatchMode.DISABLED, or use Linger(linger_ms, batch_size) for timer-based batching.
    • Call close() to flush any pending batched messages before shutdown.
    • No changes needed for consumers; batched deliveries are unpacked automatically.

Written for commit dcac3da. Summary will update on new commits.

add batch_enqueue() for explicit multi-message RPCs, smart batching via
BatchMode (AUTO/DISABLED/Linger) that routes enqueue() through a background
batcher thread, and delivery batching that unpacks ConsumeResponse.messages
repeated field. update proto to include BatchEnqueue RPC and ConsumeResponse
batched messages field. single-item optimization uses singular Enqueue RPC
to preserve error types. close() drains pending messages before disconnecting.
@vieiralucas
vieiralucas merged commit c9a83fc into mainMar 24, 2026
2 of 3 checks passed

@cubic-dev-aicubic-dev-aiBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

5 issues found across 13 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="fila/errors.py">
<violation number="1" location="fila/errors.py:75">
P2: `_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</violation>
</file>
<file name="fila/v1/service_pb2_grpc.py">
<violation number="1" location="fila/v1/service_pb2_grpc.py:8">
P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</violation>
</file>
<file name="fila/v1/messages_pb2_grpc.py">
<violation number="1" location="fila/v1/messages_pb2_grpc.py:7">
P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</violation>
</file>
<file name="fila/batcher.py">
<violation number="1" location="fila/batcher.py:102">
P1: Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</violation>
</file>
<file name="tests/test_batcher.py">
<violation number="1" location="tests/test_batcher.py:243">
P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</violation>
</file>

Reply with feedback, questions, or to request a fix. Tag @cubic-dev-ai to re-run a review.

Comment threadfila/batcher.py
return

# Pair each result with its request future.
for i, result in enumerate(resp.results):

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: Futures left permanently unresolved when server returns fewer results than batch size. If len(resp.results) < len(batch), the remaining futures never get a result or exception, causing callers to block forever on .result(). Add a fallback after the loop to fail any unmatched futures.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/batcher.py, line 102:
<comment>Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</comment>
<file context>
@@ -0,0 +1,267 @@
+ return
+
+ # Pair each result with its request future.
+ for i, result in enumerate(resp.results):
+ if i >= len(batch):
+ break
</file context>
Fix with Cubic

Comment threadfila/errors.py
code = err.code()
if code == grpc.StatusCode.NOT_FOUND:
return QueueNotFoundError(f"batch_enqueue: {err.details()}")
return RPCError(code, err.details() or "")

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: _map_batch_enqueue_error() does not return BatchEnqueueError for RPC-level batch failures, despite the new exception being documented for exactly that case.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/errors.py, line 75:
<comment>`_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</comment>
<file context>
@@ -56,3 +65,11 @@ def _map_nack_error(err: grpc.RpcError) -> FilaError:
+ code = err.code()
+ if code == grpc.StatusCode.NOT_FOUND:
+ return QueueNotFoundError(f"batch_enqueue: {err.details()}")
+ return RPCError(code, err.details() or "")
</file context>
Fix with Cubic

from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2

GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/service_pb2_grpc.py, line 8:
<comment>This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</comment>
<file context>
@@ -5,7 +5,7 @@
from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic



GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/messages_pb2_grpc.py, line 7:
<comment>This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</comment>
<file context>
@@ -4,7 +4,7 @@
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic

Comment threadtests/test_batcher.py
# Either BatchEnqueue or multiple Enqueue calls will resolve things.
for _i, f in enumerate(futures):
result = f.result(timeout=5.0)
assert result is not None

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At tests/test_batcher.py, line 243:
<comment>This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</comment>
<file context>
@@ -0,0 +1,333 @@
+ # Either BatchEnqueue or multiple Enqueue calls will resolve things.
+ for _i, f in enumerate(futures):
+ result = f.result(timeout=5.0)
+ assert result is not None
+
+ batcher.close()
</file context>
Fix with Cubic

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@vieiralucas
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

feat: python sdk batch operations & smart batching - #3

Merged
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching
Mar 24, 2026
Merged

feat: python sdk batch operations & smart batching#3
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching

Conversation

@vieiralucas

@vieiralucasvieiralucas commented Mar 24, 2026

Copy link
Copy Markdown
Member

Summary

  • Add batch_enqueue() method for explicit multi-message RPCs on both sync and async clients
  • Add smart batching via BatchMode enum (AUTO/DISABLED) and Linger dataclass, routing enqueue() through a background batcher thread by default
  • Add delivery batching: consume stream transparently unpacks ConsumeResponse.messages repeated field with backward-compatible fallback to singular message field
  • Single-item optimization: when flushing 1 message, use singular Enqueue RPC to preserve error types (QueueNotFoundError vs BatchEnqueueError)
  • close() drains pending batched messages before disconnecting
  • Update proto to include BatchEnqueue RPC and ConsumeResponse.messages repeated field
  • 16 new unit tests (batcher mechanics, flush functions, batch mode types) + 9 integration tests (batch enqueue, smart batching modes)

Test plan

  • Unit tests for _flush_single and _flush_batch with mock stubs
  • Unit tests for AutoBatcher (single message, concurrent batching, close drains, stub update)
  • Unit tests for LingerBatcher (batch size trigger, linger timeout, close drains)
  • Unit tests for BatchMode, Linger, BatchEnqueueResult types
  • Integration tests for explicit batch_enqueue() (multiple messages, single message, consume verification)
  • Integration tests for async batch_enqueue()
  • Integration tests for smart batching modes (AUTO, DISABLED, LINGER, default-is-AUTO)
  • Existing integration tests still pass (enqueue/consume/ack lifecycle, nack redelivery, API key auth)
  • ruff lint clean, mypy clean (no new errors)

Generated with Claude Code


Summary by cubic

Add explicit batch_enqueue() to the Python SDK and enable smart batching by default (AUTO) for enqueue(). Also add batched delivery in consume() and update the proto with BatchEnqueue and ConsumeResponse.messages.

  • New Features

    • Client.batch_enqueue() and AsyncClient.batch_enqueue() send many messages in one RPC; return per-message BatchEnqueueResult.
    • Smart batching via BatchMode (AUTO default, DISABLED) and Linger; background batcher with single-item optimization; close() drains pending; reconnect updates batcher stub.
    • Delivery batching: consume() transparently iterates ConsumeResponse.messages with fallback to singular message.
    • Proto changes: add BatchEnqueue RPC and ConsumeResponse.messages; new BatchEnqueueError; public exports updated.
  • Migration

    • enqueue() now routes through a background batcher by default; disable with batch_mode=BatchMode.DISABLED, or use Linger(linger_ms, batch_size) for timer-based batching.
    • Call close() to flush any pending batched messages before shutdown.
    • No changes needed for consumers; batched deliveries are unpacked automatically.

Written for commit dcac3da. Summary will update on new commits.

add batch_enqueue() for explicit multi-message RPCs, smart batching via
BatchMode (AUTO/DISABLED/Linger) that routes enqueue() through a background
batcher thread, and delivery batching that unpacks ConsumeResponse.messages
repeated field. update proto to include BatchEnqueue RPC and ConsumeResponse
batched messages field. single-item optimization uses singular Enqueue RPC
to preserve error types. close() drains pending messages before disconnecting.
@vieiralucas
vieiralucas merged commit c9a83fc into mainMar 24, 2026
2 of 3 checks passed

@cubic-dev-aicubic-dev-aiBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

5 issues found across 13 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="fila/errors.py">
<violation number="1" location="fila/errors.py:75">
P2: `_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</violation>
</file>
<file name="fila/v1/service_pb2_grpc.py">
<violation number="1" location="fila/v1/service_pb2_grpc.py:8">
P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</violation>
</file>
<file name="fila/v1/messages_pb2_grpc.py">
<violation number="1" location="fila/v1/messages_pb2_grpc.py:7">
P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</violation>
</file>
<file name="fila/batcher.py">
<violation number="1" location="fila/batcher.py:102">
P1: Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</violation>
</file>
<file name="tests/test_batcher.py">
<violation number="1" location="tests/test_batcher.py:243">
P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</violation>
</file>

Reply with feedback, questions, or to request a fix. Tag @cubic-dev-ai to re-run a review.

Comment threadfila/batcher.py
return

# Pair each result with its request future.
for i, result in enumerate(resp.results):

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: Futures left permanently unresolved when server returns fewer results than batch size. If len(resp.results) < len(batch), the remaining futures never get a result or exception, causing callers to block forever on .result(). Add a fallback after the loop to fail any unmatched futures.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/batcher.py, line 102:
<comment>Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</comment>
<file context>
@@ -0,0 +1,267 @@
+ return
+
+ # Pair each result with its request future.
+ for i, result in enumerate(resp.results):
+ if i >= len(batch):
+ break
</file context>
Fix with Cubic

Comment threadfila/errors.py
code = err.code()
if code == grpc.StatusCode.NOT_FOUND:
return QueueNotFoundError(f"batch_enqueue: {err.details()}")
return RPCError(code, err.details() or "")

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: _map_batch_enqueue_error() does not return BatchEnqueueError for RPC-level batch failures, despite the new exception being documented for exactly that case.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/errors.py, line 75:
<comment>`_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</comment>
<file context>
@@ -56,3 +65,11 @@ def _map_nack_error(err: grpc.RpcError) -> FilaError:
+ code = err.code()
+ if code == grpc.StatusCode.NOT_FOUND:
+ return QueueNotFoundError(f"batch_enqueue: {err.details()}")
+ return RPCError(code, err.details() or "")
</file context>
Fix with Cubic

from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2

GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/service_pb2_grpc.py, line 8:
<comment>This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</comment>
<file context>
@@ -5,7 +5,7 @@
from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic



GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/messages_pb2_grpc.py, line 7:
<comment>This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</comment>
<file context>
@@ -4,7 +4,7 @@
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic

Comment threadtests/test_batcher.py
# Either BatchEnqueue or multiple Enqueue calls will resolve things.
for _i, f in enumerate(futures):
result = f.result(timeout=5.0)
assert result is not None

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At tests/test_batcher.py, line 243:
<comment>This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</comment>
<file context>
@@ -0,0 +1,333 @@
+ # Either BatchEnqueue or multiple Enqueue calls will resolve things.
+ for _i, f in enumerate(futures):
+ result = f.result(timeout=5.0)
+ assert result is not None
+
+ batcher.close()
</file context>
Fix with Cubic

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@vieiralucas
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

feat: python sdk batch operations & smart batching - #3

Merged
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching
Mar 24, 2026
Merged

feat: python sdk batch operations & smart batching#3
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching

Conversation

@vieiralucas

@vieiralucasvieiralucas commented Mar 24, 2026

Copy link
Copy Markdown
Member

Summary

  • Add batch_enqueue() method for explicit multi-message RPCs on both sync and async clients
  • Add smart batching via BatchMode enum (AUTO/DISABLED) and Linger dataclass, routing enqueue() through a background batcher thread by default
  • Add delivery batching: consume stream transparently unpacks ConsumeResponse.messages repeated field with backward-compatible fallback to singular message field
  • Single-item optimization: when flushing 1 message, use singular Enqueue RPC to preserve error types (QueueNotFoundError vs BatchEnqueueError)
  • close() drains pending batched messages before disconnecting
  • Update proto to include BatchEnqueue RPC and ConsumeResponse.messages repeated field
  • 16 new unit tests (batcher mechanics, flush functions, batch mode types) + 9 integration tests (batch enqueue, smart batching modes)

Test plan

  • Unit tests for _flush_single and _flush_batch with mock stubs
  • Unit tests for AutoBatcher (single message, concurrent batching, close drains, stub update)
  • Unit tests for LingerBatcher (batch size trigger, linger timeout, close drains)
  • Unit tests for BatchMode, Linger, BatchEnqueueResult types
  • Integration tests for explicit batch_enqueue() (multiple messages, single message, consume verification)
  • Integration tests for async batch_enqueue()
  • Integration tests for smart batching modes (AUTO, DISABLED, LINGER, default-is-AUTO)
  • Existing integration tests still pass (enqueue/consume/ack lifecycle, nack redelivery, API key auth)
  • ruff lint clean, mypy clean (no new errors)

Generated with Claude Code


Summary by cubic

Add explicit batch_enqueue() to the Python SDK and enable smart batching by default (AUTO) for enqueue(). Also add batched delivery in consume() and update the proto with BatchEnqueue and ConsumeResponse.messages.

  • New Features

    • Client.batch_enqueue() and AsyncClient.batch_enqueue() send many messages in one RPC; return per-message BatchEnqueueResult.
    • Smart batching via BatchMode (AUTO default, DISABLED) and Linger; background batcher with single-item optimization; close() drains pending; reconnect updates batcher stub.
    • Delivery batching: consume() transparently iterates ConsumeResponse.messages with fallback to singular message.
    • Proto changes: add BatchEnqueue RPC and ConsumeResponse.messages; new BatchEnqueueError; public exports updated.
  • Migration

    • enqueue() now routes through a background batcher by default; disable with batch_mode=BatchMode.DISABLED, or use Linger(linger_ms, batch_size) for timer-based batching.
    • Call close() to flush any pending batched messages before shutdown.
    • No changes needed for consumers; batched deliveries are unpacked automatically.

Written for commit dcac3da. Summary will update on new commits.

add batch_enqueue() for explicit multi-message RPCs, smart batching via
BatchMode (AUTO/DISABLED/Linger) that routes enqueue() through a background
batcher thread, and delivery batching that unpacks ConsumeResponse.messages
repeated field. update proto to include BatchEnqueue RPC and ConsumeResponse
batched messages field. single-item optimization uses singular Enqueue RPC
to preserve error types. close() drains pending messages before disconnecting.
@vieiralucas
vieiralucas merged commit c9a83fc into mainMar 24, 2026
2 of 3 checks passed

@cubic-dev-aicubic-dev-aiBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

5 issues found across 13 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="fila/errors.py">
<violation number="1" location="fila/errors.py:75">
P2: `_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</violation>
</file>
<file name="fila/v1/service_pb2_grpc.py">
<violation number="1" location="fila/v1/service_pb2_grpc.py:8">
P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</violation>
</file>
<file name="fila/v1/messages_pb2_grpc.py">
<violation number="1" location="fila/v1/messages_pb2_grpc.py:7">
P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</violation>
</file>
<file name="fila/batcher.py">
<violation number="1" location="fila/batcher.py:102">
P1: Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</violation>
</file>
<file name="tests/test_batcher.py">
<violation number="1" location="tests/test_batcher.py:243">
P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</violation>
</file>

Reply with feedback, questions, or to request a fix. Tag @cubic-dev-ai to re-run a review.

Comment threadfila/batcher.py
return

# Pair each result with its request future.
for i, result in enumerate(resp.results):

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: Futures left permanently unresolved when server returns fewer results than batch size. If len(resp.results) < len(batch), the remaining futures never get a result or exception, causing callers to block forever on .result(). Add a fallback after the loop to fail any unmatched futures.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/batcher.py, line 102:
<comment>Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</comment>
<file context>
@@ -0,0 +1,267 @@
+ return
+
+ # Pair each result with its request future.
+ for i, result in enumerate(resp.results):
+ if i >= len(batch):
+ break
</file context>
Fix with Cubic

Comment threadfila/errors.py
code = err.code()
if code == grpc.StatusCode.NOT_FOUND:
return QueueNotFoundError(f"batch_enqueue: {err.details()}")
return RPCError(code, err.details() or "")

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: _map_batch_enqueue_error() does not return BatchEnqueueError for RPC-level batch failures, despite the new exception being documented for exactly that case.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/errors.py, line 75:
<comment>`_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</comment>
<file context>
@@ -56,3 +65,11 @@ def _map_nack_error(err: grpc.RpcError) -> FilaError:
+ code = err.code()
+ if code == grpc.StatusCode.NOT_FOUND:
+ return QueueNotFoundError(f"batch_enqueue: {err.details()}")
+ return RPCError(code, err.details() or "")
</file context>
Fix with Cubic

from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2

GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/service_pb2_grpc.py, line 8:
<comment>This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</comment>
<file context>
@@ -5,7 +5,7 @@
from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic



GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/messages_pb2_grpc.py, line 7:
<comment>This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</comment>
<file context>
@@ -4,7 +4,7 @@
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic

Comment threadtests/test_batcher.py
# Either BatchEnqueue or multiple Enqueue calls will resolve things.
for _i, f in enumerate(futures):
result = f.result(timeout=5.0)
assert result is not None

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At tests/test_batcher.py, line 243:
<comment>This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</comment>
<file context>
@@ -0,0 +1,333 @@
+ # Either BatchEnqueue or multiple Enqueue calls will resolve things.
+ for _i, f in enumerate(futures):
+ result = f.result(timeout=5.0)
+ assert result is not None
+
+ batcher.close()
</file context>
Fix with Cubic

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@vieiralucas
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

feat: python sdk batch operations & smart batching - #3

Merged
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching
Mar 24, 2026
Merged

feat: python sdk batch operations & smart batching#3
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching

Conversation

@vieiralucas

@vieiralucasvieiralucas commented Mar 24, 2026

Copy link
Copy Markdown
Member

Summary

  • Add batch_enqueue() method for explicit multi-message RPCs on both sync and async clients
  • Add smart batching via BatchMode enum (AUTO/DISABLED) and Linger dataclass, routing enqueue() through a background batcher thread by default
  • Add delivery batching: consume stream transparently unpacks ConsumeResponse.messages repeated field with backward-compatible fallback to singular message field
  • Single-item optimization: when flushing 1 message, use singular Enqueue RPC to preserve error types (QueueNotFoundError vs BatchEnqueueError)
  • close() drains pending batched messages before disconnecting
  • Update proto to include BatchEnqueue RPC and ConsumeResponse.messages repeated field
  • 16 new unit tests (batcher mechanics, flush functions, batch mode types) + 9 integration tests (batch enqueue, smart batching modes)

Test plan

  • Unit tests for _flush_single and _flush_batch with mock stubs
  • Unit tests for AutoBatcher (single message, concurrent batching, close drains, stub update)
  • Unit tests for LingerBatcher (batch size trigger, linger timeout, close drains)
  • Unit tests for BatchMode, Linger, BatchEnqueueResult types
  • Integration tests for explicit batch_enqueue() (multiple messages, single message, consume verification)
  • Integration tests for async batch_enqueue()
  • Integration tests for smart batching modes (AUTO, DISABLED, LINGER, default-is-AUTO)
  • Existing integration tests still pass (enqueue/consume/ack lifecycle, nack redelivery, API key auth)
  • ruff lint clean, mypy clean (no new errors)

Generated with Claude Code


Summary by cubic

Add explicit batch_enqueue() to the Python SDK and enable smart batching by default (AUTO) for enqueue(). Also add batched delivery in consume() and update the proto with BatchEnqueue and ConsumeResponse.messages.

  • New Features

    • Client.batch_enqueue() and AsyncClient.batch_enqueue() send many messages in one RPC; return per-message BatchEnqueueResult.
    • Smart batching via BatchMode (AUTO default, DISABLED) and Linger; background batcher with single-item optimization; close() drains pending; reconnect updates batcher stub.
    • Delivery batching: consume() transparently iterates ConsumeResponse.messages with fallback to singular message.
    • Proto changes: add BatchEnqueue RPC and ConsumeResponse.messages; new BatchEnqueueError; public exports updated.
  • Migration

    • enqueue() now routes through a background batcher by default; disable with batch_mode=BatchMode.DISABLED, or use Linger(linger_ms, batch_size) for timer-based batching.
    • Call close() to flush any pending batched messages before shutdown.
    • No changes needed for consumers; batched deliveries are unpacked automatically.

Written for commit dcac3da. Summary will update on new commits.

add batch_enqueue() for explicit multi-message RPCs, smart batching via
BatchMode (AUTO/DISABLED/Linger) that routes enqueue() through a background
batcher thread, and delivery batching that unpacks ConsumeResponse.messages
repeated field. update proto to include BatchEnqueue RPC and ConsumeResponse
batched messages field. single-item optimization uses singular Enqueue RPC
to preserve error types. close() drains pending messages before disconnecting.
@vieiralucas
vieiralucas merged commit c9a83fc into mainMar 24, 2026
2 of 3 checks passed

@cubic-dev-aicubic-dev-aiBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

5 issues found across 13 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="fila/errors.py">
<violation number="1" location="fila/errors.py:75">
P2: `_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</violation>
</file>
<file name="fila/v1/service_pb2_grpc.py">
<violation number="1" location="fila/v1/service_pb2_grpc.py:8">
P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</violation>
</file>
<file name="fila/v1/messages_pb2_grpc.py">
<violation number="1" location="fila/v1/messages_pb2_grpc.py:7">
P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</violation>
</file>
<file name="fila/batcher.py">
<violation number="1" location="fila/batcher.py:102">
P1: Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</violation>
</file>
<file name="tests/test_batcher.py">
<violation number="1" location="tests/test_batcher.py:243">
P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</violation>
</file>

Reply with feedback, questions, or to request a fix. Tag @cubic-dev-ai to re-run a review.

Comment threadfila/batcher.py
return

# Pair each result with its request future.
for i, result in enumerate(resp.results):

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: Futures left permanently unresolved when server returns fewer results than batch size. If len(resp.results) < len(batch), the remaining futures never get a result or exception, causing callers to block forever on .result(). Add a fallback after the loop to fail any unmatched futures.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/batcher.py, line 102:
<comment>Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</comment>
<file context>
@@ -0,0 +1,267 @@
+ return
+
+ # Pair each result with its request future.
+ for i, result in enumerate(resp.results):
+ if i >= len(batch):
+ break
</file context>
Fix with Cubic

Comment threadfila/errors.py
code = err.code()
if code == grpc.StatusCode.NOT_FOUND:
return QueueNotFoundError(f"batch_enqueue: {err.details()}")
return RPCError(code, err.details() or "")

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: _map_batch_enqueue_error() does not return BatchEnqueueError for RPC-level batch failures, despite the new exception being documented for exactly that case.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/errors.py, line 75:
<comment>`_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</comment>
<file context>
@@ -56,3 +65,11 @@ def _map_nack_error(err: grpc.RpcError) -> FilaError:
+ code = err.code()
+ if code == grpc.StatusCode.NOT_FOUND:
+ return QueueNotFoundError(f"batch_enqueue: {err.details()}")
+ return RPCError(code, err.details() or "")
</file context>
Fix with Cubic

from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2

GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/service_pb2_grpc.py, line 8:
<comment>This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</comment>
<file context>
@@ -5,7 +5,7 @@
from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic



GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/messages_pb2_grpc.py, line 7:
<comment>This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</comment>
<file context>
@@ -4,7 +4,7 @@
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic

Comment threadtests/test_batcher.py
# Either BatchEnqueue or multiple Enqueue calls will resolve things.
for _i, f in enumerate(futures):
result = f.result(timeout=5.0)
assert result is not None

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At tests/test_batcher.py, line 243:
<comment>This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</comment>
<file context>
@@ -0,0 +1,333 @@
+ # Either BatchEnqueue or multiple Enqueue calls will resolve things.
+ for _i, f in enumerate(futures):
+ result = f.result(timeout=5.0)
+ assert result is not None
+
+ batcher.close()
</file context>
Fix with Cubic

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@vieiralucas
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

feat: python sdk batch operations & smart batching - #3

Merged
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching
Mar 24, 2026
Merged

feat: python sdk batch operations & smart batching#3
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching

Conversation

@vieiralucas

@vieiralucasvieiralucas commented Mar 24, 2026

Copy link
Copy Markdown
Member

Summary

  • Add batch_enqueue() method for explicit multi-message RPCs on both sync and async clients
  • Add smart batching via BatchMode enum (AUTO/DISABLED) and Linger dataclass, routing enqueue() through a background batcher thread by default
  • Add delivery batching: consume stream transparently unpacks ConsumeResponse.messages repeated field with backward-compatible fallback to singular message field
  • Single-item optimization: when flushing 1 message, use singular Enqueue RPC to preserve error types (QueueNotFoundError vs BatchEnqueueError)
  • close() drains pending batched messages before disconnecting
  • Update proto to include BatchEnqueue RPC and ConsumeResponse.messages repeated field
  • 16 new unit tests (batcher mechanics, flush functions, batch mode types) + 9 integration tests (batch enqueue, smart batching modes)

Test plan

  • Unit tests for _flush_single and _flush_batch with mock stubs
  • Unit tests for AutoBatcher (single message, concurrent batching, close drains, stub update)
  • Unit tests for LingerBatcher (batch size trigger, linger timeout, close drains)
  • Unit tests for BatchMode, Linger, BatchEnqueueResult types
  • Integration tests for explicit batch_enqueue() (multiple messages, single message, consume verification)
  • Integration tests for async batch_enqueue()
  • Integration tests for smart batching modes (AUTO, DISABLED, LINGER, default-is-AUTO)
  • Existing integration tests still pass (enqueue/consume/ack lifecycle, nack redelivery, API key auth)
  • ruff lint clean, mypy clean (no new errors)

Generated with Claude Code


Summary by cubic

Add explicit batch_enqueue() to the Python SDK and enable smart batching by default (AUTO) for enqueue(). Also add batched delivery in consume() and update the proto with BatchEnqueue and ConsumeResponse.messages.

  • New Features

    • Client.batch_enqueue() and AsyncClient.batch_enqueue() send many messages in one RPC; return per-message BatchEnqueueResult.
    • Smart batching via BatchMode (AUTO default, DISABLED) and Linger; background batcher with single-item optimization; close() drains pending; reconnect updates batcher stub.
    • Delivery batching: consume() transparently iterates ConsumeResponse.messages with fallback to singular message.
    • Proto changes: add BatchEnqueue RPC and ConsumeResponse.messages; new BatchEnqueueError; public exports updated.
  • Migration

    • enqueue() now routes through a background batcher by default; disable with batch_mode=BatchMode.DISABLED, or use Linger(linger_ms, batch_size) for timer-based batching.
    • Call close() to flush any pending batched messages before shutdown.
    • No changes needed for consumers; batched deliveries are unpacked automatically.

Written for commit dcac3da. Summary will update on new commits.

add batch_enqueue() for explicit multi-message RPCs, smart batching via
BatchMode (AUTO/DISABLED/Linger) that routes enqueue() through a background
batcher thread, and delivery batching that unpacks ConsumeResponse.messages
repeated field. update proto to include BatchEnqueue RPC and ConsumeResponse
batched messages field. single-item optimization uses singular Enqueue RPC
to preserve error types. close() drains pending messages before disconnecting.
@vieiralucas
vieiralucas merged commit c9a83fc into mainMar 24, 2026
2 of 3 checks passed

@cubic-dev-aicubic-dev-aiBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

5 issues found across 13 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="fila/errors.py">
<violation number="1" location="fila/errors.py:75">
P2: `_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</violation>
</file>
<file name="fila/v1/service_pb2_grpc.py">
<violation number="1" location="fila/v1/service_pb2_grpc.py:8">
P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</violation>
</file>
<file name="fila/v1/messages_pb2_grpc.py">
<violation number="1" location="fila/v1/messages_pb2_grpc.py:7">
P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</violation>
</file>
<file name="fila/batcher.py">
<violation number="1" location="fila/batcher.py:102">
P1: Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</violation>
</file>
<file name="tests/test_batcher.py">
<violation number="1" location="tests/test_batcher.py:243">
P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</violation>
</file>

Reply with feedback, questions, or to request a fix. Tag @cubic-dev-ai to re-run a review.

Comment threadfila/batcher.py
return

# Pair each result with its request future.
for i, result in enumerate(resp.results):

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: Futures left permanently unresolved when server returns fewer results than batch size. If len(resp.results) < len(batch), the remaining futures never get a result or exception, causing callers to block forever on .result(). Add a fallback after the loop to fail any unmatched futures.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/batcher.py, line 102:
<comment>Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</comment>
<file context>
@@ -0,0 +1,267 @@
+ return
+
+ # Pair each result with its request future.
+ for i, result in enumerate(resp.results):
+ if i >= len(batch):
+ break
</file context>
Fix with Cubic

Comment threadfila/errors.py
code = err.code()
if code == grpc.StatusCode.NOT_FOUND:
return QueueNotFoundError(f"batch_enqueue: {err.details()}")
return RPCError(code, err.details() or "")

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: _map_batch_enqueue_error() does not return BatchEnqueueError for RPC-level batch failures, despite the new exception being documented for exactly that case.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/errors.py, line 75:
<comment>`_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</comment>
<file context>
@@ -56,3 +65,11 @@ def _map_nack_error(err: grpc.RpcError) -> FilaError:
+ code = err.code()
+ if code == grpc.StatusCode.NOT_FOUND:
+ return QueueNotFoundError(f"batch_enqueue: {err.details()}")
+ return RPCError(code, err.details() or "")
</file context>
Fix with Cubic

from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2

GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/service_pb2_grpc.py, line 8:
<comment>This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</comment>
<file context>
@@ -5,7 +5,7 @@
from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic



GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/messages_pb2_grpc.py, line 7:
<comment>This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</comment>
<file context>
@@ -4,7 +4,7 @@
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic

Comment threadtests/test_batcher.py
# Either BatchEnqueue or multiple Enqueue calls will resolve things.
for _i, f in enumerate(futures):
result = f.result(timeout=5.0)
assert result is not None

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At tests/test_batcher.py, line 243:
<comment>This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</comment>
<file context>
@@ -0,0 +1,333 @@
+ # Either BatchEnqueue or multiple Enqueue calls will resolve things.
+ for _i, f in enumerate(futures):
+ result = f.result(timeout=5.0)
+ assert result is not None
+
+ batcher.close()
</file context>
Fix with Cubic

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@vieiralucas
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

feat: python sdk batch operations & smart batching - #3

Merged
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching
Mar 24, 2026
Merged

feat: python sdk batch operations & smart batching#3
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching

Conversation

@vieiralucas

@vieiralucasvieiralucas commented Mar 24, 2026

Copy link
Copy Markdown
Member

Summary

  • Add batch_enqueue() method for explicit multi-message RPCs on both sync and async clients
  • Add smart batching via BatchMode enum (AUTO/DISABLED) and Linger dataclass, routing enqueue() through a background batcher thread by default
  • Add delivery batching: consume stream transparently unpacks ConsumeResponse.messages repeated field with backward-compatible fallback to singular message field
  • Single-item optimization: when flushing 1 message, use singular Enqueue RPC to preserve error types (QueueNotFoundError vs BatchEnqueueError)
  • close() drains pending batched messages before disconnecting
  • Update proto to include BatchEnqueue RPC and ConsumeResponse.messages repeated field
  • 16 new unit tests (batcher mechanics, flush functions, batch mode types) + 9 integration tests (batch enqueue, smart batching modes)

Test plan

  • Unit tests for _flush_single and _flush_batch with mock stubs
  • Unit tests for AutoBatcher (single message, concurrent batching, close drains, stub update)
  • Unit tests for LingerBatcher (batch size trigger, linger timeout, close drains)
  • Unit tests for BatchMode, Linger, BatchEnqueueResult types
  • Integration tests for explicit batch_enqueue() (multiple messages, single message, consume verification)
  • Integration tests for async batch_enqueue()
  • Integration tests for smart batching modes (AUTO, DISABLED, LINGER, default-is-AUTO)
  • Existing integration tests still pass (enqueue/consume/ack lifecycle, nack redelivery, API key auth)
  • ruff lint clean, mypy clean (no new errors)

Generated with Claude Code


Summary by cubic

Add explicit batch_enqueue() to the Python SDK and enable smart batching by default (AUTO) for enqueue(). Also add batched delivery in consume() and update the proto with BatchEnqueue and ConsumeResponse.messages.

  • New Features

    • Client.batch_enqueue() and AsyncClient.batch_enqueue() send many messages in one RPC; return per-message BatchEnqueueResult.
    • Smart batching via BatchMode (AUTO default, DISABLED) and Linger; background batcher with single-item optimization; close() drains pending; reconnect updates batcher stub.
    • Delivery batching: consume() transparently iterates ConsumeResponse.messages with fallback to singular message.
    • Proto changes: add BatchEnqueue RPC and ConsumeResponse.messages; new BatchEnqueueError; public exports updated.
  • Migration

    • enqueue() now routes through a background batcher by default; disable with batch_mode=BatchMode.DISABLED, or use Linger(linger_ms, batch_size) for timer-based batching.
    • Call close() to flush any pending batched messages before shutdown.
    • No changes needed for consumers; batched deliveries are unpacked automatically.

Written for commit dcac3da. Summary will update on new commits.

add batch_enqueue() for explicit multi-message RPCs, smart batching via
BatchMode (AUTO/DISABLED/Linger) that routes enqueue() through a background
batcher thread, and delivery batching that unpacks ConsumeResponse.messages
repeated field. update proto to include BatchEnqueue RPC and ConsumeResponse
batched messages field. single-item optimization uses singular Enqueue RPC
to preserve error types. close() drains pending messages before disconnecting.
@vieiralucas
vieiralucas merged commit c9a83fc into mainMar 24, 2026
2 of 3 checks passed

@cubic-dev-aicubic-dev-aiBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

5 issues found across 13 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="fila/errors.py">
<violation number="1" location="fila/errors.py:75">
P2: `_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</violation>
</file>
<file name="fila/v1/service_pb2_grpc.py">
<violation number="1" location="fila/v1/service_pb2_grpc.py:8">
P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</violation>
</file>
<file name="fila/v1/messages_pb2_grpc.py">
<violation number="1" location="fila/v1/messages_pb2_grpc.py:7">
P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</violation>
</file>
<file name="fila/batcher.py">
<violation number="1" location="fila/batcher.py:102">
P1: Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</violation>
</file>
<file name="tests/test_batcher.py">
<violation number="1" location="tests/test_batcher.py:243">
P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</violation>
</file>

Reply with feedback, questions, or to request a fix. Tag @cubic-dev-ai to re-run a review.

Comment threadfila/batcher.py
return

# Pair each result with its request future.
for i, result in enumerate(resp.results):

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: Futures left permanently unresolved when server returns fewer results than batch size. If len(resp.results) < len(batch), the remaining futures never get a result or exception, causing callers to block forever on .result(). Add a fallback after the loop to fail any unmatched futures.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/batcher.py, line 102:
<comment>Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</comment>
<file context>
@@ -0,0 +1,267 @@
+ return
+
+ # Pair each result with its request future.
+ for i, result in enumerate(resp.results):
+ if i >= len(batch):
+ break
</file context>
Fix with Cubic

Comment threadfila/errors.py
code = err.code()
if code == grpc.StatusCode.NOT_FOUND:
return QueueNotFoundError(f"batch_enqueue: {err.details()}")
return RPCError(code, err.details() or "")

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: _map_batch_enqueue_error() does not return BatchEnqueueError for RPC-level batch failures, despite the new exception being documented for exactly that case.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/errors.py, line 75:
<comment>`_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</comment>
<file context>
@@ -56,3 +65,11 @@ def _map_nack_error(err: grpc.RpcError) -> FilaError:
+ code = err.code()
+ if code == grpc.StatusCode.NOT_FOUND:
+ return QueueNotFoundError(f"batch_enqueue: {err.details()}")
+ return RPCError(code, err.details() or "")
</file context>
Fix with Cubic

from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2

GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/service_pb2_grpc.py, line 8:
<comment>This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</comment>
<file context>
@@ -5,7 +5,7 @@
from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic



GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/messages_pb2_grpc.py, line 7:
<comment>This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</comment>
<file context>
@@ -4,7 +4,7 @@
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic

Comment threadtests/test_batcher.py
# Either BatchEnqueue or multiple Enqueue calls will resolve things.
for _i, f in enumerate(futures):
result = f.result(timeout=5.0)
assert result is not None

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At tests/test_batcher.py, line 243:
<comment>This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</comment>
<file context>
@@ -0,0 +1,333 @@
+ # Either BatchEnqueue or multiple Enqueue calls will resolve things.
+ for _i, f in enumerate(futures):
+ result = f.result(timeout=5.0)
+ assert result is not None
+
+ batcher.close()
</file context>
Fix with Cubic

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@vieiralucas
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

feat: python sdk batch operations & smart batching - #3

Merged
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching
Mar 24, 2026
Merged

feat: python sdk batch operations & smart batching#3
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching

Conversation

@vieiralucas

@vieiralucasvieiralucas commented Mar 24, 2026

Copy link
Copy Markdown
Member

Summary

  • Add batch_enqueue() method for explicit multi-message RPCs on both sync and async clients
  • Add smart batching via BatchMode enum (AUTO/DISABLED) and Linger dataclass, routing enqueue() through a background batcher thread by default
  • Add delivery batching: consume stream transparently unpacks ConsumeResponse.messages repeated field with backward-compatible fallback to singular message field
  • Single-item optimization: when flushing 1 message, use singular Enqueue RPC to preserve error types (QueueNotFoundError vs BatchEnqueueError)
  • close() drains pending batched messages before disconnecting
  • Update proto to include BatchEnqueue RPC and ConsumeResponse.messages repeated field
  • 16 new unit tests (batcher mechanics, flush functions, batch mode types) + 9 integration tests (batch enqueue, smart batching modes)

Test plan

  • Unit tests for _flush_single and _flush_batch with mock stubs
  • Unit tests for AutoBatcher (single message, concurrent batching, close drains, stub update)
  • Unit tests for LingerBatcher (batch size trigger, linger timeout, close drains)
  • Unit tests for BatchMode, Linger, BatchEnqueueResult types
  • Integration tests for explicit batch_enqueue() (multiple messages, single message, consume verification)
  • Integration tests for async batch_enqueue()
  • Integration tests for smart batching modes (AUTO, DISABLED, LINGER, default-is-AUTO)
  • Existing integration tests still pass (enqueue/consume/ack lifecycle, nack redelivery, API key auth)
  • ruff lint clean, mypy clean (no new errors)

Generated with Claude Code


Summary by cubic

Add explicit batch_enqueue() to the Python SDK and enable smart batching by default (AUTO) for enqueue(). Also add batched delivery in consume() and update the proto with BatchEnqueue and ConsumeResponse.messages.

  • New Features

    • Client.batch_enqueue() and AsyncClient.batch_enqueue() send many messages in one RPC; return per-message BatchEnqueueResult.
    • Smart batching via BatchMode (AUTO default, DISABLED) and Linger; background batcher with single-item optimization; close() drains pending; reconnect updates batcher stub.
    • Delivery batching: consume() transparently iterates ConsumeResponse.messages with fallback to singular message.
    • Proto changes: add BatchEnqueue RPC and ConsumeResponse.messages; new BatchEnqueueError; public exports updated.
  • Migration

    • enqueue() now routes through a background batcher by default; disable with batch_mode=BatchMode.DISABLED, or use Linger(linger_ms, batch_size) for timer-based batching.
    • Call close() to flush any pending batched messages before shutdown.
    • No changes needed for consumers; batched deliveries are unpacked automatically.

Written for commit dcac3da. Summary will update on new commits.

add batch_enqueue() for explicit multi-message RPCs, smart batching via
BatchMode (AUTO/DISABLED/Linger) that routes enqueue() through a background
batcher thread, and delivery batching that unpacks ConsumeResponse.messages
repeated field. update proto to include BatchEnqueue RPC and ConsumeResponse
batched messages field. single-item optimization uses singular Enqueue RPC
to preserve error types. close() drains pending messages before disconnecting.
@vieiralucas
vieiralucas merged commit c9a83fc into mainMar 24, 2026
2 of 3 checks passed

@cubic-dev-aicubic-dev-aiBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

5 issues found across 13 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="fila/errors.py">
<violation number="1" location="fila/errors.py:75">
P2: `_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</violation>
</file>
<file name="fila/v1/service_pb2_grpc.py">
<violation number="1" location="fila/v1/service_pb2_grpc.py:8">
P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</violation>
</file>
<file name="fila/v1/messages_pb2_grpc.py">
<violation number="1" location="fila/v1/messages_pb2_grpc.py:7">
P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</violation>
</file>
<file name="fila/batcher.py">
<violation number="1" location="fila/batcher.py:102">
P1: Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</violation>
</file>
<file name="tests/test_batcher.py">
<violation number="1" location="tests/test_batcher.py:243">
P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</violation>
</file>

Reply with feedback, questions, or to request a fix. Tag @cubic-dev-ai to re-run a review.

Comment threadfila/batcher.py
return

# Pair each result with its request future.
for i, result in enumerate(resp.results):

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: Futures left permanently unresolved when server returns fewer results than batch size. If len(resp.results) < len(batch), the remaining futures never get a result or exception, causing callers to block forever on .result(). Add a fallback after the loop to fail any unmatched futures.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/batcher.py, line 102:
<comment>Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</comment>
<file context>
@@ -0,0 +1,267 @@
+ return
+
+ # Pair each result with its request future.
+ for i, result in enumerate(resp.results):
+ if i >= len(batch):
+ break
</file context>
Fix with Cubic

Comment threadfila/errors.py
code = err.code()
if code == grpc.StatusCode.NOT_FOUND:
return QueueNotFoundError(f"batch_enqueue: {err.details()}")
return RPCError(code, err.details() or "")

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: _map_batch_enqueue_error() does not return BatchEnqueueError for RPC-level batch failures, despite the new exception being documented for exactly that case.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/errors.py, line 75:
<comment>`_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</comment>
<file context>
@@ -56,3 +65,11 @@ def _map_nack_error(err: grpc.RpcError) -> FilaError:
+ code = err.code()
+ if code == grpc.StatusCode.NOT_FOUND:
+ return QueueNotFoundError(f"batch_enqueue: {err.details()}")
+ return RPCError(code, err.details() or "")
</file context>
Fix with Cubic

from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2

GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/service_pb2_grpc.py, line 8:
<comment>This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</comment>
<file context>
@@ -5,7 +5,7 @@
from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic



GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/messages_pb2_grpc.py, line 7:
<comment>This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</comment>
<file context>
@@ -4,7 +4,7 @@
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic

Comment threadtests/test_batcher.py
# Either BatchEnqueue or multiple Enqueue calls will resolve things.
for _i, f in enumerate(futures):
result = f.result(timeout=5.0)
assert result is not None

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At tests/test_batcher.py, line 243:
<comment>This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</comment>
<file context>
@@ -0,0 +1,333 @@
+ # Either BatchEnqueue or multiple Enqueue calls will resolve things.
+ for _i, f in enumerate(futures):
+ result = f.result(timeout=5.0)
+ assert result is not None
+
+ batcher.close()
</file context>
Fix with Cubic

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@vieiralucas
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

feat: python sdk batch operations & smart batching - #3

Merged
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching
Mar 24, 2026
Merged

feat: python sdk batch operations & smart batching#3
vieiralucas merged 1 commit into
mainfrom
feat/26.3-batch-smart-batching

Conversation

@vieiralucas

@vieiralucasvieiralucas commented Mar 24, 2026

Copy link
Copy Markdown
Member

Summary

  • Add batch_enqueue() method for explicit multi-message RPCs on both sync and async clients
  • Add smart batching via BatchMode enum (AUTO/DISABLED) and Linger dataclass, routing enqueue() through a background batcher thread by default
  • Add delivery batching: consume stream transparently unpacks ConsumeResponse.messages repeated field with backward-compatible fallback to singular message field
  • Single-item optimization: when flushing 1 message, use singular Enqueue RPC to preserve error types (QueueNotFoundError vs BatchEnqueueError)
  • close() drains pending batched messages before disconnecting
  • Update proto to include BatchEnqueue RPC and ConsumeResponse.messages repeated field
  • 16 new unit tests (batcher mechanics, flush functions, batch mode types) + 9 integration tests (batch enqueue, smart batching modes)

Test plan

  • Unit tests for _flush_single and _flush_batch with mock stubs
  • Unit tests for AutoBatcher (single message, concurrent batching, close drains, stub update)
  • Unit tests for LingerBatcher (batch size trigger, linger timeout, close drains)
  • Unit tests for BatchMode, Linger, BatchEnqueueResult types
  • Integration tests for explicit batch_enqueue() (multiple messages, single message, consume verification)
  • Integration tests for async batch_enqueue()
  • Integration tests for smart batching modes (AUTO, DISABLED, LINGER, default-is-AUTO)
  • Existing integration tests still pass (enqueue/consume/ack lifecycle, nack redelivery, API key auth)
  • ruff lint clean, mypy clean (no new errors)

Generated with Claude Code


Summary by cubic

Add explicit batch_enqueue() to the Python SDK and enable smart batching by default (AUTO) for enqueue(). Also add batched delivery in consume() and update the proto with BatchEnqueue and ConsumeResponse.messages.

  • New Features

    • Client.batch_enqueue() and AsyncClient.batch_enqueue() send many messages in one RPC; return per-message BatchEnqueueResult.
    • Smart batching via BatchMode (AUTO default, DISABLED) and Linger; background batcher with single-item optimization; close() drains pending; reconnect updates batcher stub.
    • Delivery batching: consume() transparently iterates ConsumeResponse.messages with fallback to singular message.
    • Proto changes: add BatchEnqueue RPC and ConsumeResponse.messages; new BatchEnqueueError; public exports updated.
  • Migration

    • enqueue() now routes through a background batcher by default; disable with batch_mode=BatchMode.DISABLED, or use Linger(linger_ms, batch_size) for timer-based batching.
    • Call close() to flush any pending batched messages before shutdown.
    • No changes needed for consumers; batched deliveries are unpacked automatically.

Written for commit dcac3da. Summary will update on new commits.

add batch_enqueue() for explicit multi-message RPCs, smart batching via
BatchMode (AUTO/DISABLED/Linger) that routes enqueue() through a background
batcher thread, and delivery batching that unpacks ConsumeResponse.messages
repeated field. update proto to include BatchEnqueue RPC and ConsumeResponse
batched messages field. single-item optimization uses singular Enqueue RPC
to preserve error types. close() drains pending messages before disconnecting.
@vieiralucas
vieiralucas merged commit c9a83fc into mainMar 24, 2026
2 of 3 checks passed

@cubic-dev-aicubic-dev-aiBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

5 issues found across 13 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="fila/errors.py">
<violation number="1" location="fila/errors.py:75">
P2: `_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</violation>
</file>
<file name="fila/v1/service_pb2_grpc.py">
<violation number="1" location="fila/v1/service_pb2_grpc.py:8">
P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</violation>
</file>
<file name="fila/v1/messages_pb2_grpc.py">
<violation number="1" location="fila/v1/messages_pb2_grpc.py:7">
P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</violation>
</file>
<file name="fila/batcher.py">
<violation number="1" location="fila/batcher.py:102">
P1: Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</violation>
</file>
<file name="tests/test_batcher.py">
<violation number="1" location="tests/test_batcher.py:243">
P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</violation>
</file>

Reply with feedback, questions, or to request a fix. Tag @cubic-dev-ai to re-run a review.

Comment threadfila/batcher.py
return

# Pair each result with its request future.
for i, result in enumerate(resp.results):

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: Futures left permanently unresolved when server returns fewer results than batch size. If len(resp.results) < len(batch), the remaining futures never get a result or exception, causing callers to block forever on .result(). Add a fallback after the loop to fail any unmatched futures.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/batcher.py, line 102:
<comment>Futures left permanently unresolved when server returns fewer results than batch size. If `len(resp.results) < len(batch)`, the remaining futures never get a result or exception, causing callers to block forever on `.result()`. Add a fallback after the loop to fail any unmatched futures.</comment>
<file context>
@@ -0,0 +1,267 @@
+ return
+
+ # Pair each result with its request future.
+ for i, result in enumerate(resp.results):
+ if i >= len(batch):
+ break
</file context>
Fix with Cubic

Comment threadfila/errors.py
code = err.code()
if code == grpc.StatusCode.NOT_FOUND:
return QueueNotFoundError(f"batch_enqueue: {err.details()}")
return RPCError(code, err.details() or "")

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: _map_batch_enqueue_error() does not return BatchEnqueueError for RPC-level batch failures, despite the new exception being documented for exactly that case.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/errors.py, line 75:
<comment>`_map_batch_enqueue_error()` does not return `BatchEnqueueError` for RPC-level batch failures, despite the new exception being documented for exactly that case.</comment>
<file context>
@@ -56,3 +65,11 @@ def _map_nack_error(err: grpc.RpcError) -> FilaError:
+ code = err.code()
+ if code == grpc.StatusCode.NOT_FOUND:
+ return QueueNotFoundError(f"batch_enqueue: {err.details()}")
+ return RPCError(code, err.details() or "")
</file context>
Fix with Cubic

from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2

GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/service_pb2_grpc.py, line 8:
<comment>This generated version pin raises the SDK’s runtime minimum to grpcio 1.78.1, which can break users on the declared supported range (>=1.60.0) at import time.</comment>
<file context>
@@ -5,7 +5,7 @@
from fila.v1 import service_pb2 as fila_dot_v1_dot_service__pb2
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic



GRPC_GENERATED_VERSION = '1.78.0'
GRPC_GENERATED_VERSION = '1.78.1'

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At fila/v1/messages_pb2_grpc.py, line 7:
<comment>This change raises the minimum runtime grpcio version to 1.78.1 and can break imports for users on 1.78.0 due to the generated runtime version guard.</comment>
<file context>
@@ -4,7 +4,7 @@
-GRPC_GENERATED_VERSION = '1.78.0'
+GRPC_GENERATED_VERSION = '1.78.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
</file context>
Fix with Cubic

Comment threadtests/test_batcher.py
# Either BatchEnqueue or multiple Enqueue calls will resolve things.
for _i, f in enumerate(futures):
result = f.result(timeout=5.0)
assert result is not None

@cubic-dev-aicubic-dev-aiBotMar 24, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At tests/test_batcher.py, line 243:
<comment>This assertion is too weak for a batching behavior test; it can pass even when messages are not batched.</comment>
<file context>
@@ -0,0 +1,333 @@
+ # Either BatchEnqueue or multiple Enqueue calls will resolve things.
+ for _i, f in enumerate(futures):
+ result = f.result(timeout=5.0)
+ assert result is not None
+
+ batcher.close()
</file context>
Fix with Cubic

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@vieiralucas