About TP+¶
This page introduces TP+, the single-node alternative to RT, describes its architecture, and explains how it manages deduplication, buffering, and persistence.
TP+ is the EMS implementation chosen when the stack is deployed in single-node mode, as an alternative to the multi-node RT sequencer. It isn't intended for multi-node deployments.
TP+ improves on RT in a few ways:
- Improved I/O: compared to RT, which writes separate input/output logs, TP+ writes a single set of logs, reducing disk I/O.
- Resilience: TP+ lets EMS publishers buffer outbound traffic temporarily if the TP+ instance is unavailable, giving the system time to recover instead of relying on RT's disk-first replication or suffering an outage.
- Subscriber isolation: TP+ shields itself from slow subscribers by keeping a circular buffer of recent messages and allowing lagging clients to replay directly from disk. This prevents IPC back-pressure from propagating into process-level TP+ failure (due to memory spikes) and gives administrators time to diagnose problematic consumers.
The TP+ core code keeps a circular buffer of recent messages, tracks topic maps, and only delivers live updates to subscribers that have already caught up. Subscribers that request an older position receive either a buffer replay (if within the recent range) or a log replay from disk via the bundled libtpp replay APIs.
Architecture¶
EMS delegates messaging to a set of TP+ layers:
- TP Orchestrator (TPO): loaded by the RTO service class, starts one
kxsTP_*NTASK per stream and keeps those instances registered in discovery, with relevant ancillary information, so EMS clients can resolve the stream endpoint when they publish or subscribe. - TP Core (
tpp.qandtp.q): wraps the native libtpp log writer, maintains the in-memory buffers (topics, subscribers, deduplication state), and streams writes, flushes, and checkpoints in the background. - Publisher-side library (
tppub.q): manages TCP connections to TP+ on behalf of publishers, handles deduplication group registration, implements buffering until TP+ is ready, and conforms publishes to have RT-style headers. - Subscriber-side library (
tpsub.q): manages TCP connections to TP+ on behalf of subscribers, sets up subscription to topics of interest from a defined position in the stream, and allows for replay directly from disk when subscribers are behind the in-memory buffer. Implements all subscription-related activities exposed in the EMS layer (ems.q), such as pausing, resuming, and changing topic filters.
Deduplication, buffering, and persistence¶
Each deduplication group has a disk-backed high-watermark dictionary, checkpointed regularly (driven by the checkpointFreq config in ems.yaml) with TP+, so duplicates can be dropped even after restarts. TP Core writes messages into libtpp, updates its min/current/next position window, and maintains an in-memory circular buffer containing the position, topic, and serialized message. Buffer management purges old slots using a configurable bufferSize and bufferPurgeRatio, ensuring that slow subscribers are disconnected without unbounded memory growth.
Log flush events trigger acknowledgements (when .cfgp.tpPubAck is enabled) that the publisher library consumes to clear its in-flight buffer, and the tracker updates the deduplication checkpoint so subscribers can restart from a known durable position. When publisher acknowledgements are disabled, the publisher does no buffering. The RT_LOG_PATH environment variable is set to the TP log directory, so both TP+ and RT instances share the same disk layout.