How AggregatorStream Combines Multi-Frame Features in Streaming Mode

AggregatorStream fuses multi-frame video features by running frame-level attention in parallel with causal global attention that uses a paged KV-cache, then concatenates both representations into a unified tensor of shape [B, S, P, 2C].

AggregatorStream, implemented in the Robbyant/lingbot-map repository, enables real-time video understanding by combining spatial details from individual frames with temporal context across sequences. Unlike batch processing, this streaming-capable subclass of AggregatorBase processes frames incrementally while maintaining memory of previous states through specialized caching mechanisms defined in lingbot_map/aggregator/stream.py.

Dual-Pathway Architecture for Streaming Aggregation

AggregatorStream employs two complementary attention pathways that separate local spatial processing from temporal aggregation. This architecture allows the model to handle variable-length video streams without recomputing attention over the entire history.

Frame-Level Attention Pathway

The frame-level pathway processes each video frame independently using standard Vision Transformer (ViT) blocks. In _process_frame_attention_ (lines 447-462 of lingbot_map/aggregator/stream.py), the model extracts per-frame patch tokens without cross-frame communication, preserving fine-grained spatial details from the current input.

Global Causal Attention Pathway

The global pathway implements causal cross-frame attention through _process_causal_stream_ (lines 315-329) and _process_global_attention_ (lines 360-384). This mechanism restricts each new frame to attend only to previously-seen frames, ensuring temporal consistency while preventing information leakage from future states.

Preparing Special Tokens for Variable-Length Sequences

Streaming inference requires dynamic expansion of special tokens (camera, register, scale) to match the accumulated frame history. The _prepare_special_tokens method (lines 303-368) handles this by calculating the true sequence length S_true from the KV-cache state rather than the current batch size.

The implementation uses slice_expand_and_flatten from lingbot_map/aggregator/base.py to broadcast tokens across the cached timeline:


# From lingbot_map/aggregator/stream.py

if causal_inference and S_true > S_global:
    camera_token_full = slice_expand_and_flatten(self.camera_token, B, S_true)
    camera_token = camera_token_full[-S_global:, :, :]
    # Similar expansion for register_token and scale_token

else:
    camera_token = slice_expand_and_flatten(self.camera_token, B, S_global)

This expansion ensures that special tokens align correctly with both historical cached frames and the current incoming batch before entering the attention layers.

Causal Global Attention with KV-Cache

When processing a new batch, _process_causal_stream reshapes the input tensor to [B, S_local·P, C] and manages the KV-cache through either FlashInfer's paged cache manager or an SDPA-based dictionary cache.

The method retrieves or initializes the cache manager via _get_flashinfer_manager, then iterates through global attention blocks while updating the frame counter:


# Core streaming loop from _process_causal_stream

for _ in range(self.aa_block_size):
    if self.use_sdpa:
        tokens = self.global_blocks[global_idx](
            tokens, kv_cache=self.kv_cache, ...
        )
    else:
        manager = self._get_flashinfer_manager(
            tokens.device, tokens.dtype, tokens_per_frame=P
        )
        tokens = self.global_blocks[global_idx](
            tokens, kv_cache=manager, ...
        )
    global_idx += 1
    intermediates.append(tokens.view(B, S_local, P, C))

After the first block group completes, the model updates self.total_frames_processed, which drives the special token expansion for subsequent calls. This incremental update mechanism enables constant-memory inference regardless of video length.

Fusing Local and Global Features

Following the dual-pathway processing, AggregatorStream combines outputs through channel concatenation in the forward pass (lines 602-604). This operation merges the frame-level intermediates with global causal intermediates:


# Feature fusion in the forward method

concat_inter = torch.cat(
    [frame_intermediates[i], global_intermediates[i]], 
    dim=-1
)
output_list.append(concat_inter)

The resulting tensor has shape [B, S, P, 2·C], where the first C channels encode per-frame spatial information from _process_frame_attention_ and the remaining C channels contain multi-frame temporal context gathered via the causal KV-cache. This fused representation feeds into downstream task heads for depth estimation, pose prediction, or segmentation.

Summary

  • AggregatorStream processes video through two parallel pathways: frame-level ViT blocks and causal global attention with KV-caching.
  • Special tokens expand dynamically using slice_expand_and_flatten to match cached sequence lengths calculated from S_true.
  • The _process_causal_stream method manages KV-caches via FlashInferKVCacheManager or SDPA for efficient incremental inference.
  • Final features concatenate local and global representations into a [B, S, P, 2C] tensor, combining spatial detail with temporal consistency.

Frequently Asked Questions

What distinguishes AggregatorStream from the base Aggregator class?

AggregatorStream extends AggregatorBase with causal attention mechanisms and stateful KV-caching capabilities specifically designed for streaming inference. While the base class assumes batch processing of fixed-length sequences, AggregatorStream maintains total_frames_processed and manages paged caches through _get_flashinfer_manager, enabling processing of arbitrarily long videos frame-by-frame without quadratic memory growth.

How does the KV-cache enforce causal constraints in streaming mode?

The causal constraint emerges from the attention mask implementation in _process_causal_stream, where each new frame's queries attend only to keys and values from previous positions in the cache. When using FlashInfer, the paged KV-cache manager defined in lingbot_map/layers/flashinfer_cache.py handles memory-efficient storage of historical states, while the SDPA fallback uses a dictionary-based cache structure. Both implementations ensure that position i never accesses information from positions > i.

Why does AggregatorStream concatenate rather than add the frame and global features?

Concatenation preserves the distinct characteristics of spatial and temporal representations by maintaining separate channel dimensions rather than mixing them through addition. The resulting [B, S, P, 2C] tensor explicitly separates the C channels of per-frame detail from the C channels of cross-frame context, allowing downstream heads to learn optimal weighting between local precision and temporal coherence through their own linear projections.

What hardware acceleration options does AggregatorStream support?

The implementation supports NVIDIA GPUs through FlashInfer's optimized paged attention kernels, which reduce memory overhead and accelerate causal inference for long sequences. For environments without FlashInfer support, AggregatorStream automatically falls back to PyTorch's scaled dot-product attention (SDPA) with a dictionary-based KV-cache, ensuring compatibility across different hardware configurations while maintaining the streaming capability.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →