Skip to content

feat: consume from NATS JetStream, and drop Pulsar - #15

Merged
bludot merged 1 commit into
mainfrom
feat/nats-commands
Aug 23, 2026
Merged

bludot merged 1 commit into
mainfrom
feat/nats-commands

Conversation

@bludot

@bludot bludot commented Aug 23, 2026

Copy link
Copy Markdown
Contributor

Same treatment as weeb-vip/anime-sync#22 and weeb-vip/character-staff-sync#8. Adds serve-image-sync-nats alongside serve-image-sync-kafka.

image_processor_kafka → image_processor

The _kafka suffix described the signature, not the logic: a grep for DriverMessage, RawData and Headers came back empty — it only reads Payload. So it takes the driver message as a type parameter and serves both transports:

image_processor.NewImageProcessor[*kafka.Message](store)   // kafka handler
image_processor.NewImageProcessor[*epNats.Message](store)  // nats handler

The old Pulsar-only image_processor package is deleted, which freed the accurate name.

No transform middleware — and that is correct

Unlike the CDC consumers, neither handler here has one. ep's processor hands the raw body to Event.Transform, and both drivers pass it identically:

handler(ctx, message, message.Data)   // nats driver
handler(ctx, message, msg.Value)      // kafka driver

So Payload is populated the same way on both. The CDC consumers need a transform only because their bodies carry a Debezium {schema, payload} envelope; image-sync messages are the sync services' own JSON.

NATSSTREAMNAME is empty here, on purpose

The CDC consumers bind to Debezium's ANIMEDBSTAGING stream. This one must not: image-sync is produced by the sync services, so nothing else declares a stream over it and the driver should create one from the subject. Naming Debezium's stream would graft image-sync onto a stream whose retention is sized for change events.

Pulsar removed

No Pulsar pods, no Pulsar ArgoCD apps. Checked before deleting this time — the only two pulsar mentions outside the handler were comments, not code, and image_backfill's comment describes historical data shape (two producers keyed characters differently over time), so it stays. The stale one in imagepath now says Kafka/NATS.

Note for review

internal/worker/pool.go shows a 14-line diff I did not author — gofmt -w fixed pre-existing misformatting (&Pool[T] { → &Pool[T]{, field alignment) in a file that was unformatted on main. Kept rather than reverted, since restoring badly-formatted code to keep a diff tidy is the wrong trade.

Verification

  • go build, go vet, gofmt clean; go test ./... passes
  • Docker image built as CI does; binary lists both commands:
serve-image-sync-kafka  serve-image-sync-nats

🤖 Generated with Claude Code

Adds serve-image-sync-nats alongside serve-image-sync-kafka, so staging
can move while production keeps running Kafka untouched.

image_processor_kafka becomes image_processor and takes the driver
message as a type parameter. Nothing in it reads DriverMessage, RawData
or Headers -- only Payload -- so the _kafka suffix described the
signature rather than the logic, and it now serves both transports.

No transform middleware, matching the Kafka handler. ep's processor hands
the raw body to Event.Transform and both drivers pass it the same way, so
Payload is populated identically. The CDC consumers need a transform only
because their bodies carry a Debezium envelope; these are the sync
services' own JSON.

NATSSTREAMNAME is deliberately empty here, unlike the CDC consumers:
image-sync is produced by the sync services, so no other stream declares
it and the driver should create one from the subject rather than grafting
it onto Debezium's.

Pulsar goes: no Pulsar pods or ArgoCD apps exist. internal/producer and
the pulsar processor were only reachable from the Pulsar handler.

Builds on Go 1.25, since ep v2.3.0 requires go >= 1.24.0.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@bludot
bludot merged commit bc7932a into main Aug 23, 2026
2 checks passed
@bludot

bludot commented Aug 23, 2026

Copy link
Copy Markdown
Contributor Author

🎉 This PR is included in version 1.16.0 🎉

The release is available on:

Your semantic-release bot 📦🚀

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant