Architecture
High-Level Topology
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β Client Applications β
βββββββββββββββββββββββββββββββ¬ββββββββββββββββββββββββββββββββββββ
β ConnectRPC HTTP/2 (h2c) :9000
ββββββββββββΌβββββββββββ
β API Server β
β TextToImageService β
β SpriteService β
ββββββββββββ¬βββββββββββ
β
ββββββββββββββββΌβββββββββββββββββββββββ
β β β
ββββββΌβββββ ββββββββΌβββββββ βββββββββΌββββββββ
β Databaseβ β Kafka Brokerβ β S3 / MinIO β
β (psql) β β (Redpanda) β β (sprite uploads,β
βββββββββββ ββββββββ¬βββββββ β optional) β
β subscribe βββββββββββββββββββ
ββββββββββββΌβββββββββββββββββββββββββββββββββββ
β Worker Process β
β β
β βββββββββββββββββββββββββββββββββββββββββββ β
β β TextToImage Kafka consumer β β
β β (legacy import path) β β
β βββββββββββββββββββββββββββββββββββββββββββ β
β β
β βββββββββββββββββββββββββββββββββββββββββββ β
β β GenerateImageConsumer (Kafka) β β
β ββββββββββββββββ¬βββββββββββββββββββββββββββ β
β β trigger β
β ββββββββββββββββΌβββββββββββββββββββββββββββ β
β β Temporal Worker β task queue β β
β β "image-generation" β β
β β GenerateImageWorkflow β β
β β 1. UpdateStatusActivity β gRPC β β
β β 2. GenerateImageActivity β ext API β β
β β 3. DownloadAndUploadActivity β S3 β β
β β 4. UpdateResultActivity β gRPC β β
β β 5. PublishCompletionActivity β Kafka β β
β βββββββββββββββββββββββββββββββββββββββββββββ β
β β
β βββββββββββββββββββββββββββββββββββββββββββ β
β β Ingestion Temporal Worker (optional) β β
β β task queue "harpy-eagle-image-gen" β β
β β active only when FAL_API_KEY is set β β
β β GenerateEntityImageWorkflow β β
β β fal.ai β download β MinIO β Postgres β β
β βββββββββββββββββββββββββββββββββββββββββββββ β
βββββββββββββββββββββββββββββββββββββββββββββββββ
β
ββββββββββββββββΌβββββββββββββββ
β β β
ββββββΌβββββ ββββββββΌβββββββ βββββΌβββββββββββββββββββ
βDatabase β β Kafka β β S3 / MinIO β
β(status β β (completion β β (generated images, β
β updates)β β events) β β entity images) β
βββββββββββ βββββββββββββββ βββββββββββββββββββββββββ
The service has two distinct generation paths that share infrastructure (Postgres, Kafka, Temporal, S3/MinIO) but do not share code:
- Text-to-image (
TextToImageService+imageworker) β the original, fully event-driven pipeline: API publishes to Kafka, worker consumes and drives a Temporal workflow that calls an external image API and writes status back over gRPC. - Entity image ingestion (
lib/ingestion) β a second, independently registered Temporal worker/workflow for generating and storing images for external “entities” (e.g. game characters) directly against fal.ai. It has no Kafka involvement and no dedicated gRPC service of its own; it is started only whenFAL_API_KEYis set and is invoked by starting its workflow directly against the shared Temporal client.
There is also a sprite asset catalog (SpriteService) that is unrelated
to AI generation: it’s a straightforward CRUD service, backed by Postgres,
for cataloging pre-existing RPG character-sheet sprite images (uploaded via
the CLI, not generated).
Binary Entrypoints
cmd/api/main.go β API Server
Serves TextToImageService and SpriteService over HTTP/2 h2c on :9000,
plus /health and /info (build metadata) endpoints.
Dependency wiring:
psql.Conn β TextToImageRepo β TextToImageService β gRPC handler
kafka.Conn β Producer[*GenerateImage] β TextToImageService
psql.Conn β SpriteRepo β SpriteService β gRPC handler
(optional) S3 bucket (S3_ENDPOINT) β SpriteService (reserved for a future upload RPC β unused today)
TextToImageService RPCs: ListTextToImages, GetTextToImage,
CreateTextToImage, UpdateTextToImage
On CreateTextToImage:
- Inserts record with
STATUS_STAGED - Publishes
GenerateImageproto to Kafka
SpriteService RPCs: ListSprites, GetSprite, CreateSprite (no
UpdateSprite/DeleteSprite RPC is exposed yet, though the repository
implements Update/Delete). Sprite creation is entirely CLI-driven β the CLI
uploads the sheet PNG straight to S3 itself and calls CreateSprite with an
already-known storage_url; there is no server-side upload path or
AI-generation step involved.
cmd/worker/main.go β Worker
Runs, concurrently, in a single process:
- TextToImage Kafka consumer β subscribes to the
harpyeaglev1.TextToImagetopic (groupapp-1); legacy/import path independent of the generation workflow below. - imageworker Temporal worker β task queue
image-generation; registersGenerateImageWorkflowand its activities. - GenerateImageConsumer β subscribes to the
harpyeaglev1.GenerateImageKafka topic (ML_WORKER_TOPIC, grouptemporal-worker-group); starts aGenerateImageWorkflowrun per message via the Temporal client. - ingestion Temporal worker (optional) β task queue
harpy-eagle-image-gen; registersGenerateEntityImageWorkflowand its activities. Started only whenFAL_API_KEYis set; requires a MinIO/S3 client built from theS3_*vars (mustMinIOClient, distinct from thegocloud.dev/blobbucket the imageworker/API paths use).
Layer Structure
cmd/
api/main.go β API server entrypoint, route registration
worker/main.go β worker entrypoint, wires both Temporal workers + both Kafka consumers
*.go β Cobra CLI (root, create/get/list/sprite/textimage/builtin-sprites/s3)
lib/
base/
service.go, model.go, grpc.go, consumer.go, producer.go
β generic Service[T]/Repository[T]/Producer[T] interfaces
and a generic ConnectRPC handler shared by both
TextToImage and Sprite
harpyeagle/
textimage/
service.go β TextToImageService (CRUD + Kafka publish)
consumer.go β Kafka consumer for the legacy TextToImage topic
grpc/service.go β Connect RPC handler adapter
mocks/ β mockery-generated
sprite/
service.go β SpriteService (CRUD; Create/Update bypass the
generic base service β see comments in service.go)
model.go, interface.go
grpc/service.go β Connect RPC handler adapter
imageworker/
workflow.go β GenerateImageWorkflow definition (task queue "image-generation")
activities.go β UpdateStatus, GenerateImage, DownloadAndUpload,
UpdateResult, PublishCompletion; also the fal.ai
queue-polling logic used by GenerateImageActivity
consumer.go β Kafka consumer that starts a workflow per GenerateImage message
blobstore/
blobstore.go β gocloud.dev/blob S3-compatible bucket helper
shared by the API server and imageworker
imagegen/
falai.go, generator.go β fal.ai client used by the ingestion pipeline (ImageGenerator interface)
ingestion/
workflow.go β GenerateEntityImageWorkflow (task queue "harpy-eagle-image-gen")
activities.go β BuildPrompt, GenerateImage (via imagegen.ImageGenerator),
DownloadImage, UploadToStorage (minio-go), PersistEntityImage,
SetActiveImage β writes directly to entity_images via database/sql
worker.go β registers the workflow + activities on a Temporal worker
quota/
quota.go β daily generation quota tracker (image_generation_quota table);
defined but not currently wired into any service or handler
kafka/
conn.go, consumer.go, producer.go β franz-go based Kafka client used by cmd/worker and cmd/api
repo/
conn.go β DB connection + goose migrations (embedded, run on startup)
textimage.go β texttoimages CRUD
entity_image.go β entity_images CRUD (not currently called outside ingestion/activities.go's raw SQL)
sprite.go β sprites CRUD
schemas/
harpyeagle/v1/... β generated Go code from schemas/*.proto (buf)
schemas/
harpyeagle/v1/
textimage.proto β TextToImage message + status enum
generate.proto β GenerateImage wrapper message (text-to-image only; image-to-image is defined but unimplemented)
sprite.proto β Sprite message
services/
textimage.proto β TextToImageService RPC definitions
sprite.proto β SpriteService RPC definitions
External Dependencies
| Service | Purpose | Config |
|---|---|---|
| PostgreSQL | Persistent storage (texttoimages, entity_images, sprites, image_generation_quota) | DATABASE_URL |
| Kafka / Redpanda | Job queue + completion events | KAFKA_BROKERS |
| Temporal | Workflow orchestration (two independently registered workers) | TEMPORAL_HOST |
| MinIO / S3 | Generated/uploaded image storage | S3_* vars |
| fal.ai | Image generation backend, used by both the text-to-image path (as one possible IMAGE_API_ENDPOINT) and the entity-image ingestion path | FAL_API_KEY, IMAGE_API_ENDPOINT/IMAGE_API_TOKEN |
Temporal Workflow: GenerateImageWorkflow
Task queue: image-generation
Retry policy: 3 attempts, 5s initial interval
Activity timeout: 15 minutes per activity
| Step | Activity | Description |
|---|---|---|
| 1 | UpdateStatusActivity | gRPC call to API: set status β STATUS_BUILDING |
| 2 | GenerateImageActivity | HTTP POST to IMAGE_API_ENDPOINT (fal.ai-compatible; handles both fal.ai’s async queue response and a plain synchronous JSON response) β temporary image URL |
| 3 | DownloadAndUploadActivity | Fetch image from temp URL, upload to S3/MinIO β permanent URL. If no S3 bucket is configured, falls back to the source URL unchanged. |
| 4 | UpdateResultActivity | gRPC call: set status β STATUS_COMPLETED + URL (or STATUS_FAILED + error) |
| 5 | PublishCompletionActivity | Kafka publish to completion topic; best-effort, only on success, errors logged but non-fatal |
If step 2 or 3 fails, the workflow still runs step 4 to persist STATUS_FAILED
before returning an error.
Temporal Workflow: GenerateEntityImageWorkflow (ingestion)
Task queue: harpy-eagle-image-gen
Retry policy: 3 attempts, 1s initial interval, 2x backoff, 30s max interval
| Step | Activity | Description |
|---|---|---|
| 1 | BuildPromptActivity | Assembles a prompt from entity attributes, or passes through a caller-supplied custom prompt |
| 2 | GenerateImageActivity | Calls the configured imagegen.ImageGenerator (fal.ai) |
| 3 | DownloadImageActivity | Downloads the image bytes from the provider’s temporary URL |
| 4 | UploadToStorageActivity | Uploads to S3/MinIO via minio-go |
| 5 | PersistEntityImageActivity | Inserts a row into entity_images |
| 6 | SetActiveImageActivity | Optional; if SetAsActive was requested, clears other active flags for the entity and marks this image active. Failure here is non-fatal. |
Event Flow: Text-to-Image Generation
Client β API: CreateTextToImage
β DB: INSERT status=STATUS_STAGED
β Kafka: Publish GenerateImage message (topic harpyeaglev1.GenerateImage)
Kafka β Worker: GenerateImageConsumer
β Temporal: StartWorkflow("generate-image-{uuid}")
Temporal β Worker: GenerateImageWorkflow
β gRPC β API: UpdateTextToImage(STATUS_BUILDING)
β IMAGE_API_ENDPOINT (fal.ai or compatible): POST β temp URL
β S3: upload β permanent URL (or pass through temp URL if S3 unset)
β gRPC β API: UpdateTextToImage(STATUS_COMPLETED, url) or (STATUS_FAILED, error)
β Kafka: Publish completion event (harpyeagle.texttoimages.completed) [best-effort]
Client polls: GetTextToImage(uuid) until status != STATUS_BUILDING
CLI (cmd/)
A Cobra-based CLI (built into the same binary as main.go) talks to the API
server over gRPC:
create/get/listsubcommands fortexttoimageandspritesprite create --file ... --name ...β upload a single sprite sheet PNG to S3 and register itsprite sprites-bulk --dir ...β bulk-import every.pngin a directory as a spritebuiltin-spritesβ seeds the catalog with a hard-coded roster of vendored RPG character-sheet assets (assets/builtin-sprites/*.png), matching a companion game engine’s (“goblin-shark”) built-in demo roster slot-for-slot