diff --git a/cmd/queue-operator/main.go b/cmd/queue-operator/main.go index 0795fe5..dad9d23 100644 --- a/cmd/queue-operator/main.go +++ b/cmd/queue-operator/main.go @@ -20,10 +20,10 @@ import ( "sigs.k8s.io/controller-runtime/pkg/log/zap" "sigs.k8s.io/controller-runtime/pkg/reconcile" - kmsvcv1 "forgejo.riotpiao.com/rock/kmsvc-manage/apis/kmsvc/v1" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/operator" - kmsvctemporal "forgejo.riotpiao.com/rock/kmsvc-manage/internal/temporal" + kmsvcv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/apis/kmsvc/v1" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/operator" + kmsvctemporal "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/temporal" ) func main() { diff --git a/cmd/server/main.go b/cmd/server/main.go index f604c92..013a768 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -16,17 +16,17 @@ import ( "syscall" "time" - kafkamgmtv1 "forgejo.riotpiao.com/rock/kmsvc-proto/gen/kafkamgmt/v1" + kafkamgmtv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-proto/gen/kafkamgmt/v1" "github.com/grpc-ecosystem/grpc-gateway/v2/runtime" goredis "github.com/redis/go-redis/v9" "google.golang.org/grpc" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/api/handlers" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/config" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/core/queue" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/core/reaper" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/api/handlers" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/config" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/core/queue" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/core/reaper" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) // queueDiscoveryInterval is how often the server rescans Redis for newly diff --git a/go.mod b/go.mod index 60eec29..7d6fc53 100644 --- a/go.mod +++ b/go.mod @@ -1,9 +1,9 @@ -module forgejo.riotpiao.com/rock/kmsvc-manage +module forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage go 1.26.0 require ( - forgejo.riotpiao.com/rock/kmsvc-proto v1.4.0 + forgejo.riotpiao.com/riotpiao-poimen/kmsvc-proto v1.4.0 github.com/alicebob/miniredis/v2 v2.38.0 github.com/google/uuid v1.6.0 github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 diff --git a/go.sum b/go.sum index 5d45ff2..77ded0c 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,5 @@ -forgejo.riotpiao.com/rock/kmsvc-proto v1.4.0 h1:4id+KXQhHndnlX6TC0XCYzq0MWhQWWJi09zQq4KPdR8= -forgejo.riotpiao.com/rock/kmsvc-proto v1.4.0/go.mod h1:Tvldxxalok/gCZPaUhPNLxJleRmaF0UIlLCa/MSrHg0= +forgejo.riotpiao.com/riotpiao-poimen/kmsvc-proto v1.4.0 h1:4id+KXQhHndnlX6TC0XCYzq0MWhQWWJi09zQq4KPdR8= +forgejo.riotpiao.com/riotpiao-poimen/kmsvc-proto v1.4.0/go.mod h1:Tvldxxalok/gCZPaUhPNLxJleRmaF0UIlLCa/MSrHg0= github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1Xbatp0= github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= github.com/alicebob/miniredis/v2 v2.38.0 h1:nZAzCR+Lj+Vxk4ZXzm2NuKq2O33RXj1XxJ2e2uP9jiw= diff --git a/internal/api/handlers/consumer_registry.go b/internal/api/handlers/consumer_registry.go index 418266e..fc791e1 100644 --- a/internal/api/handlers/consumer_registry.go +++ b/internal/api/handlers/consumer_registry.go @@ -5,8 +5,8 @@ import ( "fmt" "sync" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/core/queue" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/core/queue" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" ) // ConsumerRegistry lazily creates and caches one Kafka consumer per queue, diff --git a/internal/api/handlers/queue_service.go b/internal/api/handlers/queue_service.go index c408223..c761446 100644 --- a/internal/api/handlers/queue_service.go +++ b/internal/api/handlers/queue_service.go @@ -11,13 +11,13 @@ import ( "strings" "time" - kafkamgmtv1 "forgejo.riotpiao.com/rock/kmsvc-proto/gen/kafkamgmt/v1" + kafkamgmtv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-proto/gen/kafkamgmt/v1" goredis "github.com/redis/go-redis/v9" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" "google.golang.org/protobuf/types/known/timestamppb" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/core/queue" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/core/queue" ) // QueueService implements kafkamgmtv1.QueueServiceServer. diff --git a/internal/api/handlers/queue_service_test.go b/internal/api/handlers/queue_service_test.go index 7ad95c9..c517963 100644 --- a/internal/api/handlers/queue_service_test.go +++ b/internal/api/handlers/queue_service_test.go @@ -5,14 +5,14 @@ import ( "testing" "time" - kafkamgmtv1 "forgejo.riotpiao.com/rock/kmsvc-proto/gen/kafkamgmt/v1" + kafkamgmtv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-proto/gen/kafkamgmt/v1" "github.com/alicebob/miniredis/v2" goredis "github.com/redis/go-redis/v9" "github.com/twmb/franz-go/pkg/kfake" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/core/queue" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/core/queue" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) // newTestKafka starts an in-memory, wire-protocol-compatible fake Kafka diff --git a/internal/core/queue/change_visibility.go b/internal/core/queue/change_visibility.go index 20e94d4..081a0a5 100644 --- a/internal/core/queue/change_visibility.go +++ b/internal/core/queue/change_visibility.go @@ -7,7 +7,7 @@ import ( goredis "github.com/redis/go-redis/v9" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) type ChangeVisibilityService struct { diff --git a/internal/core/queue/delete_message.go b/internal/core/queue/delete_message.go index cc33579..c2e8dcf 100644 --- a/internal/core/queue/delete_message.go +++ b/internal/core/queue/delete_message.go @@ -6,8 +6,8 @@ import ( goredis "github.com/redis/go-redis/v9" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) // OffsetCommitter is the subset of kafka.Admin DeleteMessage needs to diff --git a/internal/core/queue/fifo.go b/internal/core/queue/fifo.go index e583a44..8109ebf 100644 --- a/internal/core/queue/fifo.go +++ b/internal/core/queue/fifo.go @@ -7,7 +7,7 @@ import ( goredis "github.com/redis/go-redis/v9" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) // acquireFIFOSlot claims the per-group exclusivity gate (design.md §3) so diff --git a/internal/core/queue/queue_test.go b/internal/core/queue/queue_test.go index 33d2446..5b0187d 100644 --- a/internal/core/queue/queue_test.go +++ b/internal/core/queue/queue_test.go @@ -9,8 +9,8 @@ import ( goredis "github.com/redis/go-redis/v9" "github.com/twmb/franz-go/pkg/kfake" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) // newTestKafka starts an in-memory, wire-protocol-compatible fake Kafka diff --git a/internal/core/queue/receive_message.go b/internal/core/queue/receive_message.go index 4200317..eff5b71 100644 --- a/internal/core/queue/receive_message.go +++ b/internal/core/queue/receive_message.go @@ -9,8 +9,8 @@ import ( "github.com/google/uuid" goredis "github.com/redis/go-redis/v9" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) const defaultPollInterval = 200 * time.Millisecond // design.md §2b diff --git a/internal/core/queue/send_message.go b/internal/core/queue/send_message.go index 74ec22d..7bdb37a 100644 --- a/internal/core/queue/send_message.go +++ b/internal/core/queue/send_message.go @@ -8,8 +8,8 @@ import ( "github.com/google/uuid" goredis "github.com/redis/go-redis/v9" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) // MaxMessageBodyBytes is the SQS-compatible size cap enforced at diff --git a/internal/core/queue/shard_router.go b/internal/core/queue/shard_router.go index 1902c60..773bd5c 100644 --- a/internal/core/queue/shard_router.go +++ b/internal/core/queue/shard_router.go @@ -9,8 +9,8 @@ import ( goredis "github.com/redis/go-redis/v9" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) // ShardRouter resolves a routing key to a shard by reading the cached shard diff --git a/internal/core/reaper/reaper.go b/internal/core/reaper/reaper.go index 2edbf13..9441e85 100644 --- a/internal/core/reaper/reaper.go +++ b/internal/core/reaper/reaper.go @@ -10,8 +10,8 @@ import ( goredis "github.com/redis/go-redis/v9" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/core/queue" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/core/queue" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) const ( diff --git a/internal/core/reaper/reaper_test.go b/internal/core/reaper/reaper_test.go index 3322e84..9f5e6dd 100644 --- a/internal/core/reaper/reaper_test.go +++ b/internal/core/reaper/reaper_test.go @@ -10,9 +10,9 @@ import ( goredis "github.com/redis/go-redis/v9" "github.com/twmb/franz-go/pkg/kfake" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/core/queue" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/core/queue" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) // newTestKafka mirrors internal/core/queue's test setup: an in-memory, diff --git a/internal/operator/admin.go b/internal/operator/admin.go index e0e3fd1..b61f598 100644 --- a/internal/operator/admin.go +++ b/internal/operator/admin.go @@ -3,7 +3,7 @@ package operator import ( "context" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" ) // TopicAdmin is the subset of internal/kafka.Admin the reconciler needs, diff --git a/internal/operator/fake_admin.go b/internal/operator/fake_admin.go index 57834ab..03a3a79 100644 --- a/internal/operator/fake_admin.go +++ b/internal/operator/fake_admin.go @@ -4,7 +4,7 @@ import ( "context" "sync" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" ) // fakeAdmin is an in-memory TopicAdmin for reconciler tests — avoids needing diff --git a/internal/operator/queue_controller.go b/internal/operator/queue_controller.go index 3a0935e..023c114 100644 --- a/internal/operator/queue_controller.go +++ b/internal/operator/queue_controller.go @@ -18,9 +18,9 @@ import ( "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" ctrllog "sigs.k8s.io/controller-runtime/pkg/log" - kmsvcv1 "forgejo.riotpiao.com/rock/kmsvc-manage/apis/kmsvc/v1" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + kmsvcv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/apis/kmsvc/v1" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) const finalizerName = "kmsvc.io/queue-operator" diff --git a/internal/operator/queue_controller_test.go b/internal/operator/queue_controller_test.go index 5a1a648..43e7369 100644 --- a/internal/operator/queue_controller_test.go +++ b/internal/operator/queue_controller_test.go @@ -15,9 +15,9 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client/fake" - kmsvcv1 "forgejo.riotpiao.com/rock/kmsvc-manage/apis/kmsvc/v1" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" - kmsvcredis "forgejo.riotpiao.com/rock/kmsvc-manage/internal/redis" + kmsvcv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/apis/kmsvc/v1" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" + kmsvcredis "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/redis" ) func newTestScheme(t *testing.T) *runtime.Scheme { diff --git a/internal/operator/shard_drain.go b/internal/operator/shard_drain.go index 62731fd..d1dbec5 100644 --- a/internal/operator/shard_drain.go +++ b/internal/operator/shard_drain.go @@ -5,8 +5,8 @@ import ( "fmt" "time" - kmsvcv1 "forgejo.riotpiao.com/rock/kmsvc-manage/apis/kmsvc/v1" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" + kmsvcv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/apis/kmsvc/v1" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" ) func secondsToDuration(s int32) time.Duration { diff --git a/internal/operator/shard_split.go b/internal/operator/shard_split.go index a7fe623..6cf72ff 100644 --- a/internal/operator/shard_split.go +++ b/internal/operator/shard_split.go @@ -7,8 +7,8 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - kmsvcv1 "forgejo.riotpiao.com/rock/kmsvc-manage/apis/kmsvc/v1" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" + kmsvcv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/apis/kmsvc/v1" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" ) // reconcileSplits samples each Active shard's throughput and splits any shard diff --git a/internal/operator/temporal_worker_controller.go b/internal/operator/temporal_worker_controller.go index 2cb1cfa..0ff290c 100644 --- a/internal/operator/temporal_worker_controller.go +++ b/internal/operator/temporal_worker_controller.go @@ -12,7 +12,7 @@ import ( "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" ctrllog "sigs.k8s.io/controller-runtime/pkg/log" - kmsvcv1 "forgejo.riotpiao.com/rock/kmsvc-manage/apis/kmsvc/v1" + kmsvcv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/apis/kmsvc/v1" ) const temporalWorkerFinalizerName = "kmsvc.io/temporal-worker" diff --git a/internal/operator/temporal_worker_controller_test.go b/internal/operator/temporal_worker_controller_test.go index 4015a61..831f813 100644 --- a/internal/operator/temporal_worker_controller_test.go +++ b/internal/operator/temporal_worker_controller_test.go @@ -12,7 +12,7 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client/fake" - kmsvcv1 "forgejo.riotpiao.com/rock/kmsvc-manage/apis/kmsvc/v1" + kmsvcv1 "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/apis/kmsvc/v1" ) func newTemporalWorkerTestScheme(t *testing.T) *runtime.Scheme { diff --git a/internal/redis/redis_test.go b/internal/redis/redis_test.go index ec10828..0b7e3c8 100644 --- a/internal/redis/redis_test.go +++ b/internal/redis/redis_test.go @@ -10,7 +10,7 @@ import ( "github.com/alicebob/miniredis/v2" goredis "github.com/redis/go-redis/v9" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" ) func newTestClient(t *testing.T) *goredis.Client { diff --git a/internal/redis/shardmap.go b/internal/redis/shardmap.go index 08ce409..5067042 100644 --- a/internal/redis/shardmap.go +++ b/internal/redis/shardmap.go @@ -7,7 +7,7 @@ import ( "github.com/redis/go-redis/v9" - "forgejo.riotpiao.com/rock/kmsvc-manage/internal/kafka" + "forgejo.riotpiao.com/riotpiao-poimen/kmsvc-manage/internal/kafka" ) // PutShardMap writes the active shard set for a queue (design.md §4's