fraudtwin.kafka#

Optional Kafka publication adapters.

Status: Optional

Classes#

fraudtwin.kafka.KafkaPublicationResult

Non-secret publication metadata suitable for a run manifest.

fraudtwin.kafka.KafkaPublisher

Publish mapped M24 records with deterministic pacing and acknowledgements.

fraudtwin.kafka.PublicationRecord

A deterministic, contract-encoded Kafka message.

Exceptions#

fraudtwin.kafka.KafkaConfigurationError

Raised when Kafka or Schema Registry configuration is incomplete.

fraudtwin.kafka.KafkaPublicationError

Raised when one or more Kafka records fail delivery.

Functions#

fraudtwin.kafka.publication_fingerprint

Fingerprint logical content, excluding broker-assigned transport values.

fraudtwin.kafka.publication_records

Map all observable M24 records into one deterministic publication order.

fraudtwin.kafka.publisher_from_environment

Build a publisher from environment-only connection settings.

fraudtwin.kafka.reconcile_remote_registry

Verify/register all local subjects in a Confluent-compatible registry.

fraudtwin.kafka.topic_for

Return the stable versioned topic for an M24 subject.

Constants and protocols#

Name

Reference

SUBJECTS

fraudtwin.kafka.SUBJECTS

Detailed API#

Native Kafka publication for the clean observable M24 event contracts.

exception fraudtwin.kafka.KafkaConfigurationError[source][source]

Bases: ValueError

Raised when Kafka or Schema Registry configuration is incomplete.

exception fraudtwin.kafka.KafkaPublicationError[source][source]

Bases: RuntimeError

Raised when one or more Kafka records fail delivery.

class fraudtwin.kafka.KafkaPublicationResult(topics, record_counts, mode, max_events_per_second, accelerated_time_multiplier, publication_fingerprint)[source][source]

Bases: object

Non-secret publication metadata suitable for a run manifest.

Parameters:
  • topics (dict[str, dict[str, object]])

  • record_counts (dict[str, int])

  • mode (str)

  • max_events_per_second (float | None)

  • accelerated_time_multiplier (float)

  • publication_fingerprint (str)

class fraudtwin.kafka.KafkaPublisher(*, producer, registry_client, registry, config=None, clock=time.monotonic, sleep=time.sleep)[source][source]

Bases: object

Publish mapped M24 records with deterministic pacing and acknowledgements.

Parameters:
  • producer (Any)

  • registry_client (Any)

  • registry (AvroContractRegistry)

  • config (KafkaConfig | None)

  • clock (Callable[[], float])

  • sleep (Callable[[float], None])

publish_records(records, run_id, *, mode='batch')[source][source]

Publish a pre-encoded stream without materializing all records.

Return type:

KafkaPublicationResult

Parameters:
  • records (Iterable[PublicationRecord])

  • run_id (str)

  • mode (str)

class fraudtwin.kafka.PublicationRecord(subject, topic, version, fingerprint, key, record_id, observable_time, datum, value, headers)[source][source]

Bases: object

A deterministic, contract-encoded Kafka message.

Parameters:
  • subject (str)

  • topic (str)

  • version (str)

  • fingerprint (str)

  • key (str)

  • record_id (str)

  • observable_time (datetime)

  • datum (dict[str, Any])

  • value (bytes)

  • headers (tuple[tuple[str, bytes], ...])

fraudtwin.kafka.publication_fingerprint(records)[source][source]

Fingerprint logical content, excluding broker-assigned transport values.

Return type:

str

Parameters:

records (Iterable[PublicationRecord])

fraudtwin.kafka.publication_records(behavior, run_id, *, registry=None, topic_prefix='fraudsim')[source][source]

Map all observable M24 records into one deterministic publication order.

Return type:

tuple[PublicationRecord, ...]

Parameters:
  • behavior (Any)

  • run_id (str)

  • registry (AvroContractRegistry | None)

  • topic_prefix (str)

fraudtwin.kafka.publisher_from_environment(*, registry=None, config=None)[source][source]

Build a publisher from environment-only connection settings.

Return type:

KafkaPublisher

Parameters:
  • registry (AvroContractRegistry | None)

  • config (KafkaConfig | None)

fraudtwin.kafka.reconcile_remote_registry(client, registry)[source][source]

Verify/register all local subjects in a Confluent-compatible registry.

Return type:

dict[str, int]

Parameters:
  • client (Any)

  • registry (AvroContractRegistry)

fraudtwin.kafka.topic_for(subject, prefix='fraudsim')[source][source]

Return the stable versioned topic for an M24 subject.

Return type:

str

Parameters:
  • subject (str)

  • prefix (str)