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
ConnManagermaintains an array ofConninstances (conns[]) to handle HTTP/2 push connections, withconns[0]as the primary connection.- Two independent streams (
opStreamforOperation,sqStreamforSquareEvent) separate Talk and Square data flows usingcreateAsyncReadableStream(). - Service type routing (3 for Square, 5/8 for Talk) in
_OnSignOnResponseand_OnPushResponsedirects 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
_OnPingCallbackhandles 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →