fraudtwin.kafka#
Optional Kafka publication adapters.
Status: Optional
Classes#
|
Non-secret publication metadata suitable for a run manifest. |
|
Publish mapped M24 records with deterministic pacing and acknowledgements. |
|
A deterministic, contract-encoded Kafka message. |
Exceptions#
|
Raised when Kafka or Schema Registry configuration is incomplete. |
|
Raised when one or more Kafka records fail delivery. |
Functions#
|
Fingerprint logical content, excluding broker-assigned transport values. |
|
Map all observable M24 records into one deterministic publication order. |
|
Build a publisher from environment-only connection settings. |
|
Verify/register all local subjects in a Confluent-compatible registry. |
|
Return the stable versioned topic for an M24 subject. |
Constants and protocols#
Name |
Reference |
|---|---|
|
|
Detailed API#
Native Kafka publication for the clean observable M24 event contracts.
- exception fraudtwin.kafka.KafkaConfigurationError[source][source]
Bases:
ValueErrorRaised when Kafka or Schema Registry configuration is incomplete.
- exception fraudtwin.kafka.KafkaPublicationError[source][source]
Bases:
RuntimeErrorRaised 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:
objectNon-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:
objectPublish 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])
- class fraudtwin.kafka.PublicationRecord(subject, topic, version, fingerprint, key, record_id, observable_time, datum, value, headers)[source][source]
Bases:
objectA 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)