Files
Hanzo Dev 5e3f4d808e Fix event pipeline and add admin HTTP endpoint
- Fix NATS JetStream subject mismatch: SubjectName now returns
  StreamName + ".data" to match existing stream subjects, fixing
  capture produce failures ("Unknown broker error")
- Fix OffsetFetch flexible version threshold (v6, not v8) and add
  version-aware encoding for v0-v5 (non-flexible) vs v6-v7 (flexible)
- Fix ApiVersions to use non-flexible encoding for v0-v2
- Add offset out-of-range detection in Fetch with proper error code
- Add offset validation and corruption purge in consumer offsets
- Add admin HTTP server on :9093 with /status, /topics, /groups
  endpoints for Kafka-friendly monitoring without raw NATS access
- Reduce Fetch/Heartbeat/ListOffsets request logging to DEBUG
- Set default log level to INFO
- Log RecordBatch hex on Produce for debugging
2026-02-23 16:24:18 -08:00

60 lines
1.9 KiB
Go

package main
import (
"os"
"os/signal"
"syscall"
log "github.com/hanzoai/stream/logging"
"github.com/hanzoai/stream/protocol"
"github.com/hanzoai/stream/types"
"github.com/spf13/cobra"
)
var config = types.Configuration{
PubSubUrl: "nats://localhost:4222",
BrokerHost: "localhost",
BrokerPort: 9092,
AdminPort: 9093,
NodeID: 1,
StreamReplicas: 1,
StorageType: "file",
}
func main() {
var rootCmd = &cobra.Command{
Use: "hanzo-stream",
Short: "Hanzo Stream — Kafka wire protocol gateway for Hanzo PubSub",
Run: func(cmd *cobra.Command, args []string) {
broker := protocol.NewBroker(&config)
log.SetLogLevel(log.INFO)
// Handle termination signals (e.g., Ctrl+C)
signalChannel := make(chan os.Signal, 1)
signal.Notify(signalChannel, syscall.SIGINT, syscall.SIGTERM)
go func() {
sig := <-signalChannel
log.Info("Received signal: %s. Shutting down...", sig)
broker.Shutdown()
os.Exit(0)
}()
// Start the broker
broker.Startup()
},
}
rootCmd.Flags().StringVar(&config.PubSubUrl, "pubsub-url", "nats://localhost:4222", "Hanzo PubSub server URL")
rootCmd.Flags().StringVar(&config.PubSubCredFile, "pubsub-creds", "", "Hanzo PubSub credentials file")
rootCmd.Flags().IntVar(&config.BrokerPort, "port", 9092, "Kafka listener port")
rootCmd.Flags().IntVar(&config.AdminPort, "admin-port", 9093, "Admin HTTP port (0 to disable)")
rootCmd.Flags().StringVar(&config.BrokerHost, "host", "localhost", "Advertised hostname")
rootCmd.Flags().IntVar(&config.NodeID, "node-id", 1, "Broker node ID")
rootCmd.Flags().IntVar(&config.StreamReplicas, "replicas", 1, "Hanzo Stream replica count")
rootCmd.Flags().StringVar(&config.StorageType, "storage", "file", "Hanzo Stream storage type: file or memory")
if err := rootCmd.Execute(); err != nil {
log.Panic("Failed to execute root command %v", err)
}
}