mirror of
https://github.com/warmbly/warmbly.git
synced 2026-08-19 00:01:14 +00:00
4080258606
Codec interface (Serialize / Deserialize / Name) lets payload encoding
decouple from the transport choice. Two implementations:
AvroCodec - wraps the existing kafka.Avrov2 Schema Registry client
via NewAvroFromClient. Preserves identical wire format
for production deployments already on Kafka + SR.
JSONCodec - encoding/json based. No external dependency, suitable
for self-hosters who don't want a Schema Registry.
Factory FromEnv selects via CODEC_PROVIDER (default avro).
Once codec.Codec is in the worker boot and publisher, self-hosters can
pick EVENTBUS_PROVIDER=nats CODEC_PROVIDER=json for a Schema-Registry-
free deployment.
11 tests cover JSON round-trip, nil guards, factory paths, and Avro
interface conformance.
41 lines
1.3 KiB
Go
41 lines
1.3 KiB
Go
package codec
|
|
|
|
import (
|
|
"fmt"
|
|
"os"
|
|
)
|
|
|
|
// FromEnv constructs the active Codec from environment variables.
|
|
//
|
|
// CODEC_PROVIDER=avro -> NewAvro (requires SCHEMA_REGISTRY_URL,
|
|
// SCHEMA_REGISTRY_KEY,
|
|
// SCHEMA_REGISTRY_SECRET)
|
|
// CODEC_PROVIDER=json -> NewJSON
|
|
// (unset) -> defaults to "avro" for backwards compatibility
|
|
//
|
|
// Schema Registry inputs come from env rather than parameters so callers can
|
|
// stay codec-agnostic: a self-hoster who picks json will leave them unset and
|
|
// never reach the avro branch. Operators already running the historical Kafka
|
|
// + Schema Registry stack get the same behavior as before with no config
|
|
// change.
|
|
func FromEnv() (Codec, error) {
|
|
provider := os.Getenv("CODEC_PROVIDER")
|
|
if provider == "" {
|
|
provider = "avro"
|
|
}
|
|
switch provider {
|
|
case "avro":
|
|
url := os.Getenv("SCHEMA_REGISTRY_URL")
|
|
key := os.Getenv("SCHEMA_REGISTRY_KEY")
|
|
secret := os.Getenv("SCHEMA_REGISTRY_SECRET")
|
|
if url == "" {
|
|
return nil, fmt.Errorf("codec: avro provider requires SCHEMA_REGISTRY_URL")
|
|
}
|
|
return NewAvro(url, key, secret)
|
|
case "json":
|
|
return NewJSON(), nil
|
|
default:
|
|
return nil, fmt.Errorf("codec: unknown CODEC_PROVIDER %q (want: avro, json)", provider)
|
|
}
|
|
}
|