How the LINEJS Connection Manager Handles Multiple Streams: Architecture and Implementation

The LINEJS ConnManager multiplexes Talk operations and Square events over one or more HTTP/2 push connections by maintaining two independent readable streams (opStream and sqStream) and routing incoming frames based on service type.

The LINEJS SDK relies on a sophisticated push notification system to deliver real-time Talk operations and Square events. At the heart of this system lies the ConnManager class, which orchestrates how the connection manager handles multiple streams concurrently over HTTP/2 connections. This architecture enables efficient real-time messaging while abstracting connection complexity from downstream consumers.

Managing Multiple HTTP/2 Push Connections

The ConnManager maintains an array of active connections to handle push notifications resiliently. In packages/linejs/base/push/connManager.ts, the conns: Conn[] property stores every active push connection, with conns[0] serving as the primary connection.

The initializeConn(state?, initServices?) method creates a new Conn instance and opens an HTTP/2 request to /PUSH/1/subs?.... This primary connection handles multiple logical streams through sign-on requests (buildAndSendSignOnRequest). The serviceType field in these requests determines the logical stream:

  • 3: Square fetchMyEvents
  • 5/8: Talk sync operations

This design allows a single HTTP/2 connection to carry multiplexed data for different LINE services without establishing separate TCP connections.

Dual Stream Architecture for Operation Types

Rather than mixing different event types into a single stream, ConnManager creates two distinct readable streams using createAsyncReadableStream<T>():

// packages/linejs/base/push/connManager.ts
opStream: ReadableStreamWriter<Operation>;
sqStream: ReadableStreamWriter<SquareEvent>;

this.opStream = this.createAsyncReadableStream<Operation>();
this.sqStream = this.createAsyncReadableStream<SquareEvent>();

Each stream is a custom writer that buffers chunks when consumer back-pressure is high and automatically recreates the underlying ReadableStream on close() or error(). The enqueue() method receives decoded objects from push response handlers and pushes them to the appropriate stream.

This separation ensures that Square events and Talk operations can be consumed independently without blocking each other, enabling different processing pipelines for chat messages versus community events.

Routing Incoming Frames by Service Type

When data frames arrive via _OnSignOnResponse or _OnPushResponse, the manager inspects the serviceType field and decodes the payload using the appropriate Thrift protocol. The routing logic directs each frame to its designated stream:

Service Type Decoding Path Destination Stream
3 (Square events) client.thrift.readThriftStruct → SquareEvent list sqStream.enqueue(ev)
5/8 (Talk sync) TMoreCompactProtocol or TCompactProtocol → sync_result opStream.enqueue(event)

Processing Square Events (Service Type 3)

When processing service type 3 in _OnSignOnResponse, the manager extracts the subscription ID and issues a fetchMyEvents request. Each event in the response is then individually enqueued:

for (const ev of events) {
  this.sqStream.enqueue(ev);
}

Processing Talk Operations (Service Type 5/8)

For Talk synchronization (service types 5 and 8), the manager parses the sync result using TMoreCompactProtocol or the fallback TCompactProtocol. After updating the client's sync state, each operation is streamed via opStream.enqueue().

Consuming Streams in Client Applications

The BaseClient instantiates the connection manager in its constructor (packages/linejs/base/core/mod.ts):

this.push = new ConnManager(this);

Consumers access these streams as asynchronous iterators:

// Consuming Talk operations
for await (const op of client.push.opStream.stream) {
  console.log("Talk operation:", op.type, op.revision);
}

// Consuming Square events  
for await (const ev of client.push.sqStream.stream) {
  console.log("Square event:", ev.type, ev.subtype);
}

The Polling module (packages/linejs/base/polling/mod.ts) can renew streams when starting new poll cycles:

this.client.push.sqStream.renew();

Periodic Connection Maintenance

The _OnPingCallback method handles periodic maintenance every time a ping frame arrives. When authentication tokens change, the manager closes and recreates the primary connection:

_OnPingCallback(pingId: number) {
  if (pingId % 3 === 0) {
    this.client.talk.noop().then(() => {
      const newToken = this.client.authToken;
      if (this.authToken !== newToken && newToken) {
        this.authToken = newToken;
        this.conns[0].close(); // Forces reconnection with new auth
      }
    });
  }
}

The manager also refreshes stale Square subscriptions and sends periodic InitAndRead status frames to keep HTTP/2 connections alive through network intermediaries.

Summary

  • ConnManager maintains an array of Conn instances (conns[]) to handle HTTP/2 push connections, with conns[0] as the primary connection.
  • Two independent streams (opStream for Operation, sqStream for SquareEvent) separate Talk and Square data flows using createAsyncReadableStream().
  • Service type routing (3 for Square, 5/8 for Talk) in _OnSignOnResponse and _OnPushResponse directs decoded Thrift payloads to the correct stream.
  • Back-pressure handling via custom readable stream writers ensures stable flow control when consumers process events slowly.
  • Automatic reconnection through _OnPingCallback handles auth token refreshes and connection health monitoring.

Frequently Asked Questions

How does LINEJS separate Square events from Talk operations in the same connection?

The ConnManager uses the serviceType field in sign-on and push response frames to distinguish between data types. Service type 3 triggers Square event processing via sqStream, while service types 5 and 8 route Talk operations through opStream. Both streams originate from the same HTTP/2 connection but remain logically isolated through separate ReadableStreamWriter instances.

What happens when the LINEJS authentication token expires during an active stream?

When _OnPingCallback detects an auth token change during its periodic check (every 3rd ping), it closes the primary connection using this.conns[0].close(). This forces a reconnection cycle where initializeConn establishes a fresh HTTP/2 session with the new token, while the opStream and sqStream automatically recreate their underlying readers to resume consumption without data loss.

Can consumers process Square and Talk streams simultaneously?

Yes. Because ConnManager exposes opStream and sqStream as independent asynchronous iterators, client code can consume both concurrently using for await...of loops or Promise.all() patterns. The streams operate on separate back-pressure queues, ensuring that slow processing of large Square events does not block urgent Talk operations.

Where are the Thrift protocol definitions for push frames located?

Type definitions for push frames reside in packages/linejs/base/push/connData.ts, including structures like LegyH2PushFrame and LegyH2SignOnResponseFrame. The actual Operation and SquareEvent types used by the streams are defined in packages/linejs/types/line_types.ts, which provides the TypeScript interfaces for decoded payload data.

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 →