Chever John a5b1b406ea security: remove hardcoded credentials and add envsubst for go-zero configs
- Clear default passwords in all service configs and local dev YAMLs
- Add entrypoint.sh with envsubst to resolve ${ENV} vars in go-zero YAML
- Update Dockerfiles to install gettext and use entrypoint
- Update docker-compose to pass secrets via environment and require via ${VAR:?...}
- Add .gitignore rules for .env files, add .env.example template

Co-Authored-By: Claude <noreply@anthropic.com>
2026-06-17 23:16:28 +08:00

150 lines
4.1 KiB
Go

package main
import (
"context"
"encoding/json"
"flag"
"fmt"
"log"
"os"
"os/signal"
"strings"
"syscall"
"time"
"github.com/IBM/sarama"
"github.com/neo4j/neo4j-go-driver/v5/neo4j"
)
type Config struct {
KafkaBrokers string
KafkaTopic string
KafkaGroup string
Neo4jURI string
Neo4jUser string
Neo4jPass string
}
func main() {
var cfg Config
flag.StringVar(&cfg.KafkaBrokers, "kafka-brokers", envOr("KAFKA_BROKERS", "localhost:9092"), "Kafka brokers")
flag.StringVar(&cfg.KafkaTopic, "kafka-topic", envOr("KAFKA_TOPIC", "muyu_crm.public.crm_tenant_relationship"), "CDC topic")
flag.StringVar(&cfg.KafkaGroup, "kafka-group", envOr("KAFKA_GROUP", "graphsync"), "consumer group")
flag.StringVar(&cfg.Neo4jURI, "neo4j-uri", envOr("NEO4J_URI", "bolt://localhost:7687"), "Neo4j bolt URI")
flag.StringVar(&cfg.Neo4jUser, "neo4j-user", envOr("NEO4J_USER", "neo4j"), "Neo4j user")
flag.StringVar(&cfg.Neo4jPass, "neo4j-pass", envOr("NEO4J_PASS", ""), "Neo4j password")
flag.Parse()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
driver, err := neo4j.NewDriverWithContext(cfg.Neo4jURI, neo4j.BasicAuth(cfg.Neo4jUser, cfg.Neo4jPass, ""))
if err != nil {
log.Fatalf("neo4j connect failed: %v", err)
}
defer driver.Close(ctx)
saramaCfg := sarama.NewConfig()
saramaCfg.Consumer.Group.Rebalance.GroupStrategies = []sarama.BalanceStrategy{sarama.NewBalanceStrategyRoundRobin()}
saramaCfg.Consumer.Offsets.Initial = sarama.OffsetOldest
group, err := sarama.NewConsumerGroup(strings.Split(cfg.KafkaBrokers, ","), cfg.KafkaGroup, saramaCfg)
if err != nil {
log.Fatalf("kafka consumer group failed: %v", err)
}
defer group.Close()
handler := &syncHandler{driver: driver}
go func() {
for {
if err := group.Consume(ctx, []string{cfg.KafkaTopic}, handler); err != nil {
log.Printf("consume error: %v, retrying in 5s...", err)
time.Sleep(5 * time.Second)
}
if ctx.Err() != nil {
return
}
}
}()
log.Printf("GraphSyncWorker started: brokers=%s topic=%s neo4j=%s", cfg.KafkaBrokers, cfg.KafkaTopic, cfg.Neo4jURI)
sig := make(chan os.Signal, 1)
signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
<-sig
log.Println("shutting down...")
cancel()
}
type syncHandler struct {
driver neo4j.DriverWithContext
}
func (h *syncHandler) Setup(_ sarama.ConsumerGroupSession) error { return nil }
func (h *syncHandler) Cleanup(_ sarama.ConsumerGroupSession) error { return nil }
func (h *syncHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for msg := range claim.Messages() {
if err := h.processMessage(session.Context(), msg.Value); err != nil {
log.Printf("process msg offset=%d err: %v", msg.Offset, err)
}
session.MarkMessage(msg, "")
}
return nil
}
type cdcEvent struct {
Op string `json:"op"`
After map[string]interface{} `json:"after"`
}
func (h *syncHandler) processMessage(ctx context.Context, data []byte) error {
var evt cdcEvent
if err := json.Unmarshal(data, &evt); err != nil {
return fmt.Errorf("unmarshal: %w", err)
}
if evt.After == nil {
return nil
}
session := h.driver.NewSession(ctx, neo4j.SessionConfig{AccessMode: neo4j.AccessModeWrite})
defer session.Close(ctx)
fromTenant, _ := evt.After["from_tenant_id"].(string)
toTenant, _ := evt.After["to_tenant_id"].(string)
relType, _ := evt.After["relation_type"].(string)
relID, _ := evt.After["relationship_id"].(string)
if fromTenant == "" || toTenant == "" {
return nil
}
cypher := `
MERGE (a:Tenant {tenantId: $from})
MERGE (b:Tenant {tenantId: $to})
MERGE (a)-[r:SUPPLIES {relationshipId: $relId}]->(b)
SET r.relationType = $relType, r.updatedAt = datetime()
`
_, err := session.Run(ctx, cypher, map[string]interface{}{
"from": fromTenant,
"to": toTenant,
"relId": relID,
"relType": relType,
})
if err != nil {
return fmt.Errorf("neo4j exec: %w", err)
}
log.Printf("synced: %s -> %s [%s] id=%s", fromTenant, toTenant, relType, relID)
return nil
}
func envOr(key, fallback string) string {
if v := os.Getenv(key); v != "" {
return v
}
return fallback
}