BlueTusk Sync
BlueTusk Sync materialises transaction-preserving Streams deliveries into
external destinations. It consumes BlueTusk.Streams only; it never reaches
logical-replication wire messages.
Release compatibility is specified in the public API compatibility and durable format compatibility contracts. The release-endurance guide defines the mandatory 24-hour four-connector gate.
The core pipeline owns the provisioning, snapshotting, catching-up, running, paused, rebuilding, reconciling, faulted, and stopped states. A source transaction is transformed and offered to a destination as one immutable batch. The Streams delivery is acknowledged only after the destination returns the exact commit-end position as durably handled. Duplicate delivery is safe when a destination reports the same position as already applied.
The delivery guarantee
BlueTusk guarantees that an acknowledged source transaction is never skipped. If a worker stops after writing a destination but before saving source progress, the last unconfirmed transaction is delivered again with the same source, transaction, commit-position, and change identities. Every official Sync connector provides a tested replay-safety mechanism, but the durable boundary depends on the destination. For NATS, the downstream consumer must retain the stable identity beyond JetStream’s configured deduplication window:
| Destination | Guarantee | What happens during recovery |
|---|---|---|
| PostgreSQL | Atomic state and checkpoint | The complete mutation and checkpoint commit in one PostgreSQL transaction. Recovery sees both or neither. |
| Redis | Atomic state and checkpoint | One same-slot Lua operation validates and writes the complete mutation plus checkpoint. Recovery sees both or neither. |
| OpenSearch | Replay-safe materialisation | A failed or ambiguous bulk is replayed with the same external versions. Older source versions cannot replace newer materialized state, and the checkpoint advances only after every item succeeds. |
| NATS JetStream | Durable publication with stable identity | JetStream acknowledges durable storage and deduplicates the deterministic message ID inside its configured window. The same ID remains in the envelope so consumers can enforce a longer boundary. |
In practical terms, PostgreSQL and Redis provide one atomic destination effect; OpenSearch converges to one version-protected materialized state after replay; and NATS provides one durable transaction identity with a clearly bounded broker-deduplication window. Third-party connectors remain conservatively at-least-once unless their own destination contract proves a stronger outcome.
SyncDestinationConformanceSuite verifies changed-payload redelivery both in
the same process and after a new connector instance starts. The official
connector tests then verify destination-specific atomicity, partial-write,
deduplication-window, versioning, and checkpoint failure boundaries. The public
SyncDestinationCapabilities flags expose which mechanics a connector supports.
Transform definitions carry a canonical SHA-256 fingerprint. A mismatch moves
the pipeline to Rebuilding and requires an explicit rebuild or migration;
BlueTusk does not silently reinterpret existing destination data.
CompositeSyncTransform builds that fingerprint from the source mapper and
every ordered transform stage. SyncPredicateTransformStage supplies explicit,
versioned filtering. JsonSyncTransformStage supplies bounded JSON redaction,
root enrichment, object flattening, and tenant routing. Its configuration is
canonicalised into the pipeline fingerprint, so changing a predicate version,
redaction path, enrichment value, flatten separator, tenant path, or size limit
requires an explicit rebuild or migration.
JSON materialisation accepts object-valued application/json payloads only and
enforces input, output, and configuration bounds. Redaction paths always use
dotted source paths even when flattened output uses another separator. Tenant
routing is resolved before redaction, rejects missing or empty tenant values by
default, requires deletes to retain their partition key, and rejects unscoped
collection deletes. Stages may filter and rewrite mapped content but cannot
invent a transaction change ID or snapshot row ID. Those rules keep replay and
destination deduplication stable.
Transformation sandbox
SandboxedSyncTransformStage runs untrusted transformation configuration as a
finite JSON instruction program. It does not load assemblies, compile
expressions, invoke delegates, use reflection, or expose file, network, process,
clock, environment, or service-provider access. The only operations are remove,
set, copy, route, conditional drop, and equality requirement. There are no
branches other than those finite predicates and no loops.
var sandbox = new SandboxedSyncTransformStage(
new SyncTransformSandboxOptions
{
Name = "public-orders",
Version = "v4",
Instructions =
[
SyncSandboxInstruction.RequireEquals("kind", "\"order\""),
SyncSandboxInstruction.Copy("customer.name", "displayName"),
SyncSandboxInstruction.Remove("customer.email"),
SyncSandboxInstruction.Set("metadata.source", "\"cdc\""),
SyncSandboxInstruction.Route("tenant.id"),
SyncSandboxInstruction.DropWhenEquals("status", "\"cancelled\""),
],
MaximumDocumentBytes = 256 * 1024,
MaximumBatchBytes = 8 * 1024 * 1024,
MaximumOperationsPerBatch = 250_000,
MaximumExecutionTime = TimeSpan.FromSeconds(2),
});
Construction canonicalises literal JSON and includes instruction order and every execution limit in the transform fingerprint. Changing the program or a limit therefore requires the same explicit destination rebuild or migration as any other transform-version change.
Execution is bounded independently by instruction count, mutation count, input and output document size, aggregate batch bytes, JSON depth, operation count, cancellation, and elapsed time. A violation becomes a poison record and follows the pipeline’s configured pause/quarantine policy. Stable change IDs and snapshot row IDs are preserved. A routing program rejects collection-wide deletes and, by default, requires ordinary deletes to retain the partition key established by the source mapper.
This is the V1 sandbox contract. It deliberately does not claim that an in-process C# delegate or a child process without an operating-system policy is sandboxed. Native, managed-assembly, and general-purpose script plug-ins remain outside the trusted boundary.
Reproduce the sandbox contract suite with:
dotnet test tests/BlueTusk.Sync.Tests/BlueTusk.Sync.Tests.csproj --filter FullyQualifiedName~SyncTransformTests
Poison transformations pause by default. QuarantineAndAdvance is accepted only
with an explicit durable quarantine sink, and the source delivery is not
acknowledged until that sink confirms storage. Destination outages and rejected
durability confirmations nack the delivery and fault the pipeline for safe
redelivery.
Quarantine replay
QuarantineAndPause is the replay-oriented poison policy. It durably stores the
quarantine record, acknowledges the source transaction, and leaves the pipeline
paused at that exact boundary. SyncQuarantineReplayCoordinator then reads the
original transaction from an ISyncQuarantineReplaySource, verifies the stored,
requested, and running transform fingerprints, reruns the transform, applies the
transaction, and compare-and-set resolves the quarantine record. Destination
application always precedes resolution, so a crash between them safely retries
an idempotent destination operation.
PostgreSqlRelaySyncQuarantineReplaySource reads the exact retained transaction
from the durable relay without moving any consumer-group checkpoint. If relay
retention has expired, replay returns SourceTransactionUnavailable and leaves
the record unresolved. PostgreSQL, Redis, OpenSearch, and NATS implement the
replay destination contract. PostgreSQL, Redis, and OpenSearch also implement
the durable read/resolve quarantine store; NATS can use any separately
configured store.
Materialising destinations atomically reject replay when their checkpoint is
already beyond the quarantined commit position. This prevents an old upsert or
delete from overwriting later state; the operator must use authoritative
reconciliation or rebuild instead. Unscoped collection deletes are never
replayed into a materialised destination. QuarantineAndAdvance remains an
explicit continue-on-poison policy, but its record may therefore become
non-replayable after later destination progress.
The dashboard exposes the separately authorised and audited
ReplayQuarantine operation when a pipeline reports quarantined transactions.
The application operation handler selects the unresolved record and passes the
request operation ID into the replay coordinator, preserving idempotency across
HTTP retries. Resume the worker only after replay reports Completed or
AlreadyCompleted.
PostgreSQL, NATS JetStream, Redis, OpenSearch, signed webhook, and Kafka
connector slices are implemented and pass the same executable
snapshot-plus-stream recovery contract. The shared count, key-set, and
partitioned content-hash engine plus
PostgreSQL, Redis, and OpenSearch repair paths are implemented and live-tested.
The destination-neutral rebuild coordinator is implemented with an explicit
cutover barrier, and the in-process hosting package provides named workers,
health, telemetry, and one-way worker handoff. The production relay
change-stream adapter provides bounded reads, continuous fenced lease renewal,
and acknowledgement-after-destination ordering. Restart-aware relay snapshot
bootstrap now reserves retention before export, restarts abandoned epochs, and
resumes completed pipelines without resetting the destination. Dashboard
read models now expose pipeline state, sampled throughput, checkpoint lag,
failures, quarantine, retries, and rebuild/reconciliation state without leaking
worker exception messages. Separately authorized retry, reconcile, and rebuild
controls now pass exact confirmations through the audit-before-mutation control
plane executor. The coordinated 1.1.0-rc.1 package train is public. Stable
1.1.0 remains non-publishable until the exact 24-hour endurance and final
stable-release gates pass.
The executable Sync release endurance gate repeatedly runs the core/hosting contracts and all seven destination suites for 24 hours, then emits a versioned evidence report. A local one-cycle smoke proves the orchestration only; it does not satisfy the release gate.
Retry, rate limits, and backpressure
Sync never starts another source transaction while the current delivery is in destination work, retry delay, or rate-limit delay. This sequential pull model propagates backpressure to Streams without an unbounded pipeline-owned queue and preserves the source transaction order.
SyncRetryOptions applies a bounded exponential backoff with configurable
jitter and a hard attempt ceiling. No exception is retried by default: an
application must register ISyncRetryClassifier, or its destination must
implement that interface, and explicitly classify each failure as transient.
The transform runs once and every attempt receives the same immutable batch or
quarantine record, including stable IDs and timestamps. Retry exhaustion faults
the pipeline and nacks the active Streams delivery; the checkpoint cannot move
past an unconfirmed destination position.
SyncRateLimitOptions can independently cap source transactions per second and
transformed bytes per second. Each destination attempt, including a retry,
consumes capacity. Waiting occurs inline before the attempt, so rate limiting
does not introduce reordering or a second buffering layer. Snapshot batches use
the transformed-byte limit but do not count as CDC transactions.
Runtime status exposes cumulative retry attempts and throttle duration.
BlueTusk.Sync.DependencyInjection publishes them through
bluetusk.sync.retries and bluetusk.sync.throttle.duration; the health
registry includes the same values for control-plane consumers. Retry classifiers
are resolved from dependency injection and remain operator policy rather than a
connector silently guessing whether a database, broker, or HTTP failure is safe
to retry.
Shared destination conformance
BlueTusk.Sync.Testing contains SyncDestinationConformanceSuite, the single
connector acceptance scenario used by PostgreSQL, NATS JetStream, Redis,
OpenSearch, Kafka, and signed webhook destinations. Service-backed suites run
against their real connector; deterministic transports cover protocol and
failure boundaries that would otherwise be timing-dependent. It verifies:
- provisioning and transform-fingerprint ownership;
- idempotent snapshot batches and completed snapshot state after a new destination instance starts;
- exact durable commit positions, same-instance duplicate delivery, and process-restart redelivery without replacing accepted content;
- explicit
RebuildRequiredresults for transform-version drift; and - durable, idempotent quarantine for connectors that expose
ISyncQuarantineSink.
An in-memory reference harness also proves the kit rejects a destination that reports a checkpoint beyond the applied source transaction. The live variants run in their existing connector CI jobs; PostgreSQL runs across versions 15–19.
Reconciliation and repair
SyncReconciler supports three explicit depths: count, partitioned key set, and
partitioned exact-content SHA-256. Key partitions are derived from the high
32 bits of SHA-256 and compared in bounded streams ordered by hash and UTF-8 key.
Results retain a configurable number of representative differences while exact
totals continue to accumulate. Count-only equality is deliberately reported as
count equality; it does not claim content equality.
Repair is unavailable in count-only mode. For key-set or content-hash runs, the
authoritative reader must include replacement content and the destination must
implement ISyncRepairSink. Repairs are idempotent upserts/deletes sent in
bounded batches. A repaired result remains a mismatch and sets
RequiresVerification; only a subsequent clean comparison proves convergence.
Repair never advances the source transaction checkpoint.
SyncPipeline.ReconcileAsync serializes reconciliation with delivery, exposes
the Reconciling state, restores the previous running/paused state after
success, and faults with diagnostics on a reader or repair failure.
Zero-downtime rebuilds
SyncRebuildCoordinator creates or resumes an isolated destination generation,
runs a consistent snapshot with full-epoch restart after exporter loss, and
then acquires an ISyncRebuildCutoverBarrier. The barrier must quiesce the
active pipeline and capture the durable relay head as an exact transaction
commit-end position. It remains held while the new generation consumes
transaction-preserving catch-up, verifies, and performs the destination’s atomic
routing swap. This makes the cutover target an observed post-snapshot boundary
instead of an unsafe position guessed before the snapshot began.
Verification has two mandatory layers. ISyncRebuildDestination first checks
generation-owned metadata and storage integrity. ISyncRebuildVerifier then
checks the rebuilding materialisation against the authoritative transformed
source view. SyncReconciliationRebuildVerifier supplies the default bounded
implementation and accepts only non-repairing partitioned content-hash requests;
count equality or old-generation equality cannot authorise a cutover because a
new transform may legitimately change both content and cardinality.
Every catch-up transaction is acknowledged only after the destination confirms the exact commit-end position. Wrong positions, transform failures, stream reordering, or an early finite stream nack the active delivery where possible, release the barrier, and leave the rebuilding generation inactive. Verification failure also releases the barrier without activation. Progress reports are observational and cannot alter durability semantics.
After the routing swap, the coordinator explicitly commits the worker handoff while the barrier is still held. Disposing an uncommitted lease resumes the previous worker; once handoff begins, disposal must keep that worker quiesced so it cannot process the old transform against the activated generation. A handoff failure reports completed activation and requires forward operator recovery, never rollback.
Previous-generation retirement is optional and occurs only after activation.
If retirement fails, SyncRebuildRetirementException explicitly reports that
activation completed and must not be rolled back. Running the coordinator again
is safe: destination preparation resumes its durable build metadata, while an
interrupted snapshot starts a new epoch and clears the isolated generation
before replay.
In-process hosting and cutover
BlueTusk.Sync.DependencyInjection registers named pipelines with
AddBlueTuskSync().AddHostedPipeline<TTransform, TDestination>(). Each worker
owns one long-lived dependency-injection scope, provisions its destination,
runs the Streams snapshot-then-stream coordinator, and leaves a faulted worker
isolated while the other registered pipelines continue. Source factories remain
explicit so an application cannot accidentally share a replication session or
consumer group between pipelines.
AddHostedPipelineSource<TTransform, TDestination>() accepts an
ISyncPipelineSource when the source owns a restart-aware lifecycle.
PostgreSqlRelaySyncPipelineSource acquires and renews its independently
checkpointed group before asking PostgreSQL to export a snapshot. A versioned
relay snapshot run records the transform fingerprint, epoch, and consistent
LSN. Reserved state left by exporter/session loss causes a safe new epoch;
completed state resumes the group without another destination reset.
After snapshot completion, retained relay transactions through the snapshot’s
consistent LSN are acknowledged as already represented. Later transactions are
passed to SyncPipeline, which advances the group only after the destination
confirms the exact commit-end LSN. Nack, abandonment, lease loss, or a process
crash therefore causes safe at-least-once redelivery. A Latest group still
cannot prove a no-gap bootstrap and is not used by this source’s default
earliest-retained reservation protocol.
BlueTuskSyncHealthRegistry exposes immutable operational snapshots. The
readiness check is unhealthy for faults, transform rebuild requirements, or
hosting errors; paused/stopped-only deployments are degraded. Activities and
metrics use the stable BlueTusk.Sync instrumentation name and report bounded
pipeline identifiers, acknowledged transaction counts, snapshot rows, failures,
and transaction duration.
AddRebuildCutover<TPositionProvider, THandoffHandler>() connects the shared
rebuild coordinator to a hosted active worker. The barrier waits for the current
delivery boundary, verifies the worker is running, captures a durable target,
and holds all further consumption through verification and activation. A failed
or cancelled pre-activation rebuild releases the worker. Once
CompleteHandoffAsync begins, the old worker is permanently cancelled before
the restart-safe handoff handler is invoked and can never resume against the new
generation.
AddPostgreSqlRelayRebuildCutover<THandoffHandler>() uses the production relay’s
separate control data source through
PostgreSqlRelaySyncCutoverPositionProvider. The target is the latest durable
relay commit position, or the snapshot’s consistent baseline when no later
transaction exists. The rebuild source and handoff handler must use an
independent relay group whose checkpoint is bound to that same snapshot epoch.
BlueTusk.Sync.Aspire wires source, destination, and—by default—a distinct
durable-relay control resource into an Aspire worker. Its options carry the
pipeline, group, transform version, destination protocol, reconciliation, and
rebuild settings through standard hierarchical configuration. Direct-slot mode
is an explicit helper that omits the control resource; the durable-relay helper
rejects using the source database itself as control storage.
PostgreSQL destination
BlueTusk.Sync.PostgreSql stores an opaque materialised document collection and
the pipeline checkpoint in the same PostgreSQL database transaction. It locks
the pipeline row, skips mutation work for an already-applied commit position,
and advances the checkpoint only after every mutation succeeds. A custom
IPostgreSqlSyncMutationWriter can target application-specific tables while
retaining the same atomic checkpoint boundary.
The default writer folds repeated operations to the final per-key result while preserving collection-delete ordering, then sends bounded multi-row commands instead of one database round trip per document. The live acceptance suite covers batches beyond one command chunk as well as retry deduplication.
The default document writer also exposes server-partitioned hash reconciliation and transactional repair. PostgreSQL computes the shared SHA-256 partition in SQL, streams rows in deterministic order, and applies a bounded repair batch in one database transaction without touching the pipeline checkpoint. A custom mutation writer does not advertise reconciliation because BlueTusk cannot infer how to inspect or repair an application-owned schema.
Snapshot reset, batches, and completion are guarded by the active snapshot epoch and transform fingerprint. The destination also implements a durable, deduplicated quarantine sink. Document and transaction byte ceilings are validated before opening the write transaction.
NATS JetStream destination
BlueTusk.Sync.Nats publishes one versioned binary envelope for each source
transaction. It waits for JetStream’s persistence acknowledgement before
returning the exact durable source position, so the Sync pipeline cannot
acknowledge a partially published transaction. Snapshot reset, start, batch,
and completion are also individually durable envelopes.
Every publish uses a fixed-size SHA-256 message ID derived from the pipeline, source, transform version, and transaction or snapshot identity. JetStream deduplicates redelivery inside its configured duplicate window; the same stable identity remains inside the envelope so downstream consumers can deduplicate beyond that window. BlueTusk still advertises at-least-once delivery.
The envelope has a magic header, explicit format version, bounded payload size,
and SHA-256 integrity footer. Consumers decode it with
NatsSyncEnvelopeReader. Mutation records retain stable change or snapshot row
IDs, collection/key routing, content type, partition key, and opaque content.
Provisioning creates a file-backed, limits-retained JetStream stream by default.
The stream carries ownership metadata for the envelope format, pipeline, source,
transform, and subject. Existing metadata and retention settings are validated
before publishing; drift pauses provisioning. A transform fingerprint change
returns RebuildRequired, so operators must provision a new stream generation
or explicitly migrate routing rather than reinterpret existing events.
The duplicate window must cover the expected worker recovery interval. Retain
stable IDs downstream even when using a long window because redelivery after
the window is valid at-least-once behaviour. Set CreateStream to false when
stream creation is managed externally; BlueTusk will still validate the stream
contract.
For local acceptance, start a JetStream-enabled NATS server and run:
$env:BLUETUSK_NATS_URL = 'nats://localhost:4222'
dotnet test tests/BlueTusk.Sync.Nats.Tests/BlueTusk.Sync.Nats.Tests.csproj
The live suite proves whole-transaction persistence, duplicate recovery after a destination restart, snapshot lifecycle deduplication, transform-generation rejection, and stable stream message counts.
Apache Kafka destination
BlueTusk.Sync.Kafka publishes one versioned JSON envelope for each complete
PostgreSQL source transaction. The event record and the destination checkpoint
are committed together with a Confluent Kafka transactional producer. The
checkpoint lives in a separate compacted state topic and is loaded with
isolation.level=read_committed on every provision or restart. BlueTusk does
not acknowledge Streams until Kafka confirms that atomic broker transaction.
This is an explicit end-to-end at-least-once contract. Kafka’s producer transaction prevents a visible event without its BlueTusk checkpoint, while stable delivery and mutation IDs let consumers make their own business effect idempotent. It does not claim that an arbitrary downstream consumer’s database write is magically part of the Kafka transaction. If the broker outcome is ambiguous, the destination is invalidated and the checkpoint is not advanced; re-provisioning reloads the authoritative compacted state before retrying.
Each pipeline owns exactly two topics:
<prefix>.events, containing ordered transaction and snapshot envelopes;<prefix>.state, withcleanup.policy=compact, containing configuration, the highest committed PostgreSQL position, and bounded current-snapshot progress.
Both topics require exactly one partition per pipeline. This is deliberate:
transaction order is a correctness boundary, not a throughput hint. Scale by
using independent pipeline/topic prefixes. A production cluster should use a
replication factor of three, acks=all, idempotence, TLS/SASL credentials from
secret-backed Confluent configuration, and a unique transactional ID for each
active pipeline writer.
var destination = new KafkaSyncDestination(new KafkaSyncOptions
{
BootstrapServers = configuration["Kafka:BootstrapServers"]!,
TopicPrefix = "bluetusk.orders",
TransactionalId = "orders-sync-primary",
ClientId = "orders-sync",
ReplicationFactor = 3,
ClientConfiguration = new Dictionary<string, string>
{
["security.protocol"] = "SaslSsl",
["sasl.mechanism"] = "SCRAM-SHA-512",
["sasl.username"] = configuration["Kafka:Username"]!,
["sasl.password"] = configuration["Kafka:Password"]!,
},
});
Provisioning validates topic existence, partition count, and compaction on the
state topic. The state record owns the pipeline, PostgreSQL source fingerprint,
and transform fingerprint. Reusing a topic prefix for another pipeline/source
fails closed; a transform change returns RebuildRequired. Snapshot batch
progress is monotonic per table, old epoch keys are tombstoned on reset, and
the ordinary CDC checkpoint never moves during snapshot delivery.
S3 and Parquet lake destination
BlueTusk.Sync.S3 stores each transaction or snapshot batch as one immutable,
Zstandard-compressed Parquet object. Every row retains the stable source ID,
mutation order, operation, collection/key routing, content type, partition key,
and opaque content bytes. File metadata records the format, delivery, event,
pipeline, and transform identities. Mutation count and final compressed byte
size are bounded before any checkpoint can advance.
S3 has no multi-object transaction, so BlueTusk uses an explicit data-lake commit protocol:
- write the deterministic Parquet data key with
If-None-Match: *; - verify an existing object has the same SHA-256 if a retry races it;
- write an immutable JSON manifest with
If-None-Match: *last; and - acknowledge the exact PostgreSQL position only after the manifest succeeds.
Lake readers enumerate commits/, never data/. A crash between steps 1 and 3
can leave an orphaned data object, but it cannot create a visible commit or
advance Streams. Retrying uses the same object keys and hashes. A conflicting
object fails closed. Transaction manifest names begin with the fixed-width
commit-end LSN so a reader can retain authoritative PostgreSQL order.
This is at-least-once delivery with immutable commit markers, not a claim that
S3 provides cross-object transactions. The marker is the durable duplicate
receipt. Re-delivery after restart returns AlreadyApplied when that marker is
present, even if the caller presents changed content for the same source
identity.
var destination = new S3SyncDestination(new S3SyncOptions
{
Client = amazonS3,
BucketName = "company-data-lake",
Prefix = "bluetusk/orders/v1",
ServerSideEncryption = ServerSideEncryptionMethod.AWSKMS,
KmsKeyId = configuration["DataLakeKmsKeyId"],
MaxMutationCount = 100_000,
MaxParquetBytes = 64 * 1024 * 1024,
});
The prefix owns one immutable configuration object containing the pipeline,
PostgreSQL source, and transform fingerprints. Reusing it for another source
fails; transform drift returns RebuildRequired. Use a new generation prefix
for rebuilds. Production defaults request S3-managed encryption, while KMS can
be enforced explicitly. Bucket policy should limit the writer to its prefix,
deny unencrypted transport, retain commit manifests, and apply lifecycle rules
to known orphan/data retention separately from committed evidence.
Signed webhook destination
BlueTusk.Sync.Webhooks delivers one bounded JSON envelope for an entire source
transaction. It also emits explicit snapshot reset, start, batch, and completion
events. Each request has a deterministic delivery ID and an HMAC-SHA256
signature over the timestamp and exact request body. HTTPS is mandatory unless
insecure HTTP is explicitly enabled for a local test.
The receiver—not an in-process cache—owns durable duplicate detection. It must
persist the delivery ID with the application result before replying with
BlueTusk-Delivery-Status: applied. A redelivery replies with duplicate.
BlueTusk treats a successful HTTP response without either acknowledgement as
ambiguous and does not advance the Streams checkpoint. HTTP 408, 425, 429, and
5xx responses receive bounded retries with the same delivery ID, body,
timestamp, and signature.
This is a concrete at-least-once contract: a crash can cause the same request to arrive again, but the receiver has a stable ID with which to suppress repeated work. It is not described as “exactly once.” Transaction envelopes are ordered per destination instance and contain stable mutation IDs, collection/key routing, content types, partition keys, and base64-encoded opaque content.
Provisioning is also signed. The receiver returns its durable transform
fingerprint in BlueTusk-Transform-Fingerprint; a different fingerprint causes
RebuildRequired before data delivery begins. Configure the destination with a
dedicated HttpClient, an HTTPS endpoint, a non-secret key identifier, and at
least 32 random signing-key bytes:
var destination = new WebhookSyncDestination(new WebhookSyncOptions
{
Client = httpClient,
Endpoint = new Uri("https://receiver.example.com/bluetusk/sync"),
KeyId = "orders-2026-09",
SigningKey = Convert.FromBase64String(configuration["WebhookSigningKey"]!),
});
The receiver must reject stale timestamps according to its clock-skew policy,
look up the key by BlueTusk-Key-Id, compare signatures in constant time,
validate formatVersion, and persist the delivery ID in the same durability
boundary as the business effect. Rotate keys by accepting both identifiers for
the maximum retry interval before removing the old key.
Redis destination
BlueTusk.Sync.Redis stores materialised documents and the source checkpoint in
one Redis Lua operation. All keys for a pipeline use the same generated Redis
Cluster hash tag, so an atomic batch never crosses slots. The script checks the
source, transform, monotonic fixed-width commit position, key types, and every
operation before writing; a predictable failure therefore cannot leave a
partial transaction or advance its checkpoint.
Repeated mutations are folded to their final per-key outcome before the script runs, while the last collection delete remains ordered before subsequent upserts. Configurable document, transaction-byte, and mutation-count ceilings bound Lua execution time and Redis argument memory. Transactions beyond those limits pause safely for operator action instead of blocking Redis indefinitely.
Documents use a small versioned binary value with the stable source change or
snapshot-row ID, content type, partition key, opaque content, and a SHA-256
integrity footer. Applications can inspect a materialised value with
ReadDocumentAsync or decode an exported value with
RedisSyncDocumentReader.
Snapshot reset atomically removes registered materialisations and clears the
checkpoint before activating a new epoch. Snapshot batches are idempotent, and
completion prevents late batches for the epoch. Quarantine records use a stable
transaction field and HSET NX, so retrying quarantine-and-advance cannot add
duplicates.
Quarantine values use a versioned JSON document and compare-and-set resolution; legacy newline records remain readable. Replay uses one Lua script to verify the source, transform, and exact checkpoint/transaction boundary, apply all mutations, and advance that checkpoint atomically.
Redis format version 2 maintains a same-slot sorted reconciliation index beside each collection hash. Lua writes update the document, index, registry, and CDC checkpoint atomically. Partition reads use bounded score ranges instead of rescanning the whole collection, and repair updates the document hash and index in one Lua call without changing the CDC checkpoint.
The live Redis suite deliberately introduces a wrong-type destination key and proves preflight rejection occurs before any mutation. It also covers retry, restart, collection-delete ordering, snapshot reset/completion, quarantine, and transform rebuild requirements.
OpenSearch destination
BlueTusk.Sync.OpenSearch uses one bounded NDJSON bulk request as the source
transaction delivery unit. It assigns SHA-256 document IDs and PostgreSQL
commit-end LSNs as external_gte versions, so a partial request or ambiguous
network failure can safely replay the entire transaction. The per-generation
checkpoint is written only after every bulk item succeeds. A partially accepted
bulk therefore never advances the checkpoint, and its successful items remain
idempotent on retry.
OpenSearch bulk operations are independently applied by the server, so this
connector deliberately does not advertise TransactionalBatches or a
co-located checkpoint. Transaction preservation comes from bounded whole-batch
submission, item-by-item response validation, stable external versions, and
checkpoint-after-bulk ordering. Collection resets complete before the
subsequent folded mutations are sent. JSON objects are the only accepted
materialisation content.
Format version 2 creates a generation-owned reconciliation sidecar beside every materialised index. Each sidecar record contains the original logical key, its shared unsigned SHA-256 partition hash, the exact content hash, content type, and routing value; application JSON remains untouched and hashed document IDs never need to be reversed. A source mutation and its sidecar operation share the same replay-safe external version in one bulk request. Partial bulk failure cannot advance the checkpoint, and replay heals either half before progress is claimed. Count reads reject materialised/sidecar cardinality drift instead of silently comparing an incomplete view.
Partitioned sidecar scans use bounded search_after pages ordered by key hash
and logical key. Repair looks up prior routing, removes an old routed copy when
the routing value changes, and writes the application document plus sidecar
without changing the CDC checkpoint. A subsequent reconciliation run is still
required to prove convergence. Logical keys and page sizes have explicit
operator-configured ceilings so reconciliation cannot create unbounded terms or
responses.
Each transform generation writes to isolated concrete indexes. Stable aliases
are attached to the initial generation, while a rebuild generation remains
invisible. BeginRebuildAsync is restart-safe and copies the active collection
registry, snapshot and catch-up writes target only the rebuild indexes,
VerifyRebuildAsync checks every active/rebuild count, and
CompleteRebuildAsync moves all aliases in one atomic OpenSearch aliases
request. Previous generations are retained until an explicit
RetireGenerationAsync call.
OpenSearch also implements ISyncRebuildDestination, allowing the shared
coordinator to drive those connector-native generation operations without
depending on OpenSearch types.
The control index owns versioned pipeline, collection, snapshot, checkpoint,
and quarantine documents. Source and transform fingerprints are validated on
every restart; a changed transform returns RebuildRequired. Index names,
aliases, and document IDs contain hashes instead of application keys, and
document, mutation-count, and encoded bulk-byte limits are checked before
submission.
Quarantine resolution uses OpenSearch sequence-number/primary-term compare and set. Replay shares the normal external-versioned bulk path and rejects an already-advanced checkpoint before applying an old transaction.
For local acceptance, run an OpenSearch node without the security plug-in and then execute:
$env:BLUETUSK_OPENSEARCH_URL = 'http://localhost:9200'
dotnet test tests/BlueTusk.Sync.OpenSearch.Tests/BlueTusk.Sync.OpenSearch.Tests.csproj
The CI and local live suite use OpenSearch 3.7.0. It deliberately causes a mapping conflict after another item has succeeded, repairs and replays the same transaction, and then covers checkpoint deduplication, collection reset, snapshot lifecycle, quarantine, restart, transform isolation, count verification, atomic alias cutover, old-generation retirement, paged content-hash reconciliation, bounded repair, and checkpoint non-advancement. The design follows the official Bulk API and Manage Aliases API contracts.