seaweedfs/weed/mq/broker
Chris Lu 02773a6107
Some checks are pending
go: build dev binaries / cleanup (push) Waiting to run
go: build dev binaries / build_dev_linux_windows (amd64, linux) (push) Blocked by required conditions
go: build dev binaries / build_dev_linux_windows (amd64, windows) (push) Blocked by required conditions
go: build dev binaries / build_dev_darwin (amd64, darwin) (push) Blocked by required conditions
go: build dev binaries / build_dev_darwin (arm64, darwin) (push) Blocked by required conditions
docker: build dev containers / build-dev-containers (push) Waiting to run
End to End / FUSE Mount (push) Waiting to run
go: build binary / Build (push) Waiting to run
Ceph S3 tests / Ceph S3 tests (push) Waiting to run
Accumulated changes for message queue (#6600)
* rename

* set agent address

* refactor

* add agent sub

* pub messages

* grpc new client

* can publish records via agent

* send init message with session id

* fmt

* check cancelled request while waiting

* use sessionId

* handle possible nil stream

* subscriber process messages

* separate debug port

* use atomic int64

* less logs

* minor

* skip io.EOF

* rename

* remove unused

* use saved offsets

* do not reuse session, since always session id is new after restart

remove last active ts from SessionEntry

* simplify printing

* purge unused

* just proxy the subscription, skipping the session step

* adjust offset types

* subscribe offset type and possible value

* start after the known tsns

* avoid wrongly set startPosition

* move

* remove

* refactor

* typo

* fix

* fix changed path
2025-03-09 23:49:42 -07:00
..
broker_connect.go adjust errors 2024-05-16 11:02:48 -07:00
broker_grpc_admin.go subscriber find broker leader first 2024-02-05 23:14:25 -08:00
broker_grpc_assign.go Add message queue agent (#6463) 2025-01-20 22:19:27 -08:00
broker_grpc_balance.go rename 2024-05-21 10:02:07 -07:00
broker_grpc_configure.go Add message queue agent (#6463) 2025-01-20 22:19:27 -08:00
broker_grpc_lookup.go Add message queue agent (#6463) 2025-01-20 22:19:27 -08:00
broker_grpc_pub_balancer.go rename 2024-05-21 10:02:07 -07:00
broker_grpc_pub_follow.go merge current message queue code changes (#6201) 2024-11-04 12:08:25 -08:00
broker_grpc_pub.go adds locking 2024-08-11 13:06:01 -07:00
broker_grpc_sub_coordinator.go balance subscribers 2024-05-27 17:30:16 -07:00
broker_grpc_sub_follow.go merge current message queue code changes (#6201) 2024-11-04 12:08:25 -08:00
broker_grpc_sub.go Accumulated changes for message queue (#6600) 2025-03-09 23:49:42 -07:00
broker_grpc_topic_partition_control.go Merge accumulated changes related to message queue (#5098) 2023-12-11 12:05:54 -08:00
broker_server.go sort imports 2024-06-14 11:44:42 -07:00
broker_topic_conf_read_write.go merge current message queue code changes (#6201) 2024-11-04 12:08:25 -08:00
broker_topic_partition_read_write.go merge current message queue code changes (#6201) 2024-11-04 12:08:25 -08:00
broker_write.go Added tls for http clients (#5766) 2024-07-16 23:14:09 -07:00