Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
37 commits
Select commit Hold shift + click to select a range
eebff88
fix: performance issues with Kafka caused by encapsulating the MQ int…
withchao Jul 28, 2025
031b37f
Merge branch 'openimsdk:main' into main
withchao Jul 28, 2025
55ebedc
Merge branch 'openimsdk:main' into main
withchao Jul 29, 2025
d9c3504
fix: admin token in standalone mode
withchao Aug 1, 2025
fd3c07f
Merge branch 'openimsdk:main' into main
withchao Aug 1, 2025
9c20271
Merge branch 'openimsdk:main' into main
withchao Oct 15, 2025
9a1d2a8
fix: full id version
withchao Oct 15, 2025
30fd83b
Merge branch 'openimsdk:main' into main
withchao Dec 12, 2025
ebda95f
fix: resolve deadlock in cache eviction and improve GetBatch implemen…
withchao Dec 12, 2025
83de399
Merge branch 'openimsdk:main' into main
withchao Dec 19, 2025
a1dd79a
refactor: replace LongConn with ClientConn interface and simplify mes…
withchao Dec 19, 2025
9da7db2
refactor: replace LongConn with ClientConn interface and simplify mes…
withchao Dec 19, 2025
6e90026
Merge branch 'openimsdk:main' into main
withchao Dec 25, 2025
79aae5a
Merge branch 'openimsdk:main' into main
withchao Dec 31, 2025
8e27646
Merge branch 'openimsdk:main' into main
withchao Jan 15, 2026
c27d331
fix: seq use $setOnInsert for min_seq in conversation update
withchao Jan 15, 2026
ec33d6c
Merge branch 'openimsdk:main' into main
withchao Jan 22, 2026
5d451fa
feat: add error code for handled friend requests and improve error ha…
withchao Jan 22, 2026
7709b75
Merge branch 'openimsdk:main' into main
withchao Jan 23, 2026
75ebca4
Merge branch 'openimsdk:main' into main
withchao Jun 5, 2026
82f8755
refactor(msg): update regex pattern for conversationID to include a t…
withchao Jun 5, 2026
e769084
Merge branch 'openimsdk:main' into main
withchao Jun 25, 2026
b5ccf80
Merge branch 'openimsdk:main' into main
withchao Jun 30, 2026
c9c7db0
Merge branch 'openimsdk:main' into main
withchao Jul 3, 2026
d0366c4
refactor: streamline cache initialization and enhance queue engine co…
withchao Jul 3, 2026
d172cfe
refactor: streamline cache initialization and enhance queue engine co…
withchao Jul 3, 2026
ad2735a
feat: add RegisterIP to API config and implement Redis server registr…
withchao Jul 3, 2026
ee67269
feat(redis): implement standalone gateway registration with Redis
withchao Jul 3, 2026
26e22ef
chore: update openimsdk/tools dependency to v0.0.50-alpha.121
withchao Jul 3, 2026
fa411a7
fix: change timer to ticker for improved periodic execution
withchao Jul 3, 2026
bfe8bcd
refactor: replace Kafka topic configuration with mqbuild constants fo…
withchao Jul 3, 2026
5882597
feat(redis): implement Redis-based locking mechanism for cron tasks
withchao Jul 3, 2026
c603db1
Merge branch 'openimsdk:main' into main
withchao Jul 6, 2026
fc3e18f
Merge branch 'openimsdk:main' into main
withchao Jul 7, 2026
10baf8f
feat: Support Streaming Messages in the Open-Source Server
withchao Jul 13, 2026
0e5a798
Merge branch 'openimsdk:main' into main
withchao Aug 4, 2026
7fb7d55
chore: update openimsdk/tools dependency to v0.0.50-alpha.122
withchao Aug 4, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,8 @@ require (
github.com/gorilla/websocket v1.5.1
github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0
github.com/mitchellh/mapstructure v1.5.0
github.com/openimsdk/protocol v0.0.73-alpha.19
github.com/openimsdk/tools v0.0.50-alpha.121
github.com/openimsdk/protocol v0.0.73-alpha.20
github.com/openimsdk/tools v0.0.50-alpha.122
github.com/pkg/errors v0.9.1 // indirect
github.com/prometheus/client_golang v1.18.0
github.com/stretchr/testify v1.11.1
Expand Down
8 changes: 4 additions & 4 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -361,10 +361,10 @@ github.com/onsi/gomega v1.25.0 h1:Vw7br2PCDYijJHSfBOWhov+8cAnUf8MfMaIOV323l6Y=
github.com/onsi/gomega v1.25.0/go.mod h1:r+zV744Re+DiYCIPRlYOTxn0YkOLcAnW8k1xXdMPGhM=
github.com/openimsdk/gomake v0.0.17 h1:q8haP48VOH45WhJRiLj1YSBJyUFJqD8CTedH65i1YH8=
github.com/openimsdk/gomake v0.0.17/go.mod h1:nnjS8yCtrPJAt1knMbyPiUwCH2gpyBzj/EZAONfUOXg=
github.com/openimsdk/protocol v0.0.73-alpha.19 h1:CvXoDF2U73UcMhLnrtMFks2Aw+bXiDgH8AITEt783/s=
github.com/openimsdk/protocol v0.0.73-alpha.19/go.mod h1:WF7EuE55vQvpyUAzDXcqg+B+446xQyEba0X35lTINmw=
github.com/openimsdk/tools v0.0.50-alpha.121 h1:TXKKgtkeMeqIs0vpolbW8rIEngE9xlESq+0NV+FoLH0=
github.com/openimsdk/tools v0.0.50-alpha.121/go.mod h1:I0WESSa7ghPIo9BL+ETlH/qEIbO6+KZioM1jwNuDwz0=
github.com/openimsdk/protocol v0.0.73-alpha.20 h1:9MnACSi6IKv2iqlxHYUJG9mgt9gyPRHxE7Lq8dxAoxI=
github.com/openimsdk/protocol v0.0.73-alpha.20/go.mod h1:WF7EuE55vQvpyUAzDXcqg+B+446xQyEba0X35lTINmw=
github.com/openimsdk/tools v0.0.50-alpha.122 h1:bKP6hrJ6kGyGUqTGuyqUqPdZxY7mZSsh86aRQJ/6Jt8=
github.com/openimsdk/tools v0.0.50-alpha.122/go.mod h1:I0WESSa7ghPIo9BL+ETlH/qEIbO6+KZioM1jwNuDwz0=
github.com/pelletier/go-toml/v2 v2.2.2 h1:aYUidT7k73Pcl9nb2gScu7NSrKCSHIDE89b3+6Wq+LM=
github.com/pelletier/go-toml/v2 v2.2.2/go.mod h1:1t835xjRzz80PqgE6HHgN2JOsmgYu/h4qDAS4n929Rs=
github.com/pierrec/lz4/v4 v4.1.21 h1:yOVMLb6qSIDP67pl/5F7RepeKYu/VmTyEXvuMI5d9mQ=
Expand Down
7 changes: 5 additions & 2 deletions internal/api/msg.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,12 +79,13 @@ func getMsgDataDescriptor() []protoreflect.FieldDescriptor {
type MessageApi struct {
Client msg.MsgClient
userClient *rpcli.UserClient
authClient *rpcli.AuthClient
imAdminUserID []string
validate *validator.Validate
}

func NewMessageApi(client msg.MsgClient, userClient *rpcli.UserClient, imAdminUserID []string) MessageApi {
return MessageApi{Client: client, userClient: userClient, imAdminUserID: imAdminUserID, validate: validator.New()}
func NewMessageApi(client msg.MsgClient, userClient *rpcli.UserClient, authClient *rpcli.AuthClient, imAdminUserID []string) MessageApi {
return MessageApi{Client: client, userClient: userClient, authClient: authClient, imAdminUserID: imAdminUserID, validate: validator.New()}
}

func (*MessageApi) SetOptions(options map[string]bool, value bool) {
Expand Down Expand Up @@ -219,6 +220,8 @@ func (m *MessageApi) getSendMsgReq(c *gin.Context, req apistruct.SendMsg) (sendM
data = &apistruct.CustomElem{}
case constant.MarkdownText:
data = &apistruct.MarkdownTextElem{}
case constant.Stream:
data = &apistruct.StreamMsgElem{}
case constant.Quote:
data = &apistruct.QuoteElem{}
case constant.OANotification:
Expand Down
5 changes: 4 additions & 1 deletion internal/api/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -253,7 +253,7 @@ func newGinRouter(ctx context.Context, client discovery.SvcDiscoveryRegistry, cf
objectGroup.GET("/*name", t.ObjectRedirect)
}
// Message
m := NewMessageApi(msg.NewMsgClient(msgConn), rpcli.NewUserClient(userConn), cfg.Share.IMAdminUser.UserIDs)
m := NewMessageApi(msg.NewMsgClient(msgConn), rpcli.NewUserClient(userConn), rpcli.NewAuthClient(authConn), cfg.Share.IMAdminUser.UserIDs)
{
msgGroup := r.Group("/msg")
msgGroup.POST("/newest_seq", m.GetSeq)
Expand All @@ -277,6 +277,9 @@ func newGinRouter(ctx context.Context, client discovery.SvcDiscoveryRegistry, cf
msgGroup.POST("/send_simple_msg", m.SendSimpleMessage)
msgGroup.POST("/check_msg_is_send_success", m.CheckMsgIsSendSuccess)
msgGroup.POST("/get_server_time", m.GetServerTime)
msgGroup.POST("/get_stream_msg", m.GetStreamMsg)
msgGroup.POST("/append_stream_msg", m.AppendStreamMsg)
msgGroup.PUT("/append_stream_msg", m.PutStreamMsg)
}
// Conversation
{
Expand Down
216 changes: 216 additions & 0 deletions internal/api/stream_msg.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,216 @@
package api

import (
"bufio"
"bytes"
"context"
"fmt"
"io"
"net/http"
"time"
"unicode/utf8"

"github.com/gin-gonic/gin"

"github.com/openimsdk/protocol/constant"
"github.com/openimsdk/protocol/msg"
"github.com/openimsdk/tools/a2r"
"github.com/openimsdk/tools/apiresp"
"github.com/openimsdk/tools/errs"
"github.com/openimsdk/tools/log"
)

func (m *MessageApi) GetStreamMsg(c *gin.Context) {
a2r.Call(c, msg.MsgClient.GetStreamMsg, m.Client)
}

func (m *MessageApi) AppendStreamMsg(c *gin.Context) {
a2r.Call(c, msg.MsgClient.AppendStreamMsg, m.Client)
}

func (m *MessageApi) PutStreamMsg(c *gin.Context) {
var (
conversationID string
clientMsgID string
)
{
operationID := c.GetHeader(constant.OperationID)
if operationID == "" {
operationID = c.Query(constant.OperationID)
}
if operationID == "" {
m.putErr(c, errs.ErrArgs.WrapMsg("operationID is empty"))
return
}
c.Set(constant.OperationID, operationID)
conversationID = c.Query("conversationID")
if conversationID == "" {
conversationID = c.GetHeader("conversationID")
}
if conversationID == "" {
m.putErr(c, errs.ErrArgs.WrapMsg("conversationID is empty"))
return
}
clientMsgID = c.Query("clientMsgID")
if clientMsgID == "" {
clientMsgID = c.GetHeader("clientMsgID")
}
if clientMsgID == "" {
m.putErr(c, errs.ErrArgs.WrapMsg("clientMsgID is empty"))
return
}
token := c.GetHeader("token")
if token == "" {
token = c.Query("token")
}
if token == "" {
m.putErr(c, errs.ErrTokenInvalid.WrapMsg("token is empty"))
return
}
resp, err := m.authClient.ParseToken(c, token)
if err != nil {
m.putErr(c, err)
return
}
c.Set(constant.OpUserPlatform, constant.PlatformIDToName(int(resp.PlatformID)))
c.Set(constant.OpUserID, resp.UserID)
}
done := make(chan struct{})
streamCh := make(chan string, 8)

go func() {
defer func() {
close(streamCh)
c.Request.Body.Close()
}()
buf := make([]byte, 256)
body := NewUTF8Reader(c.Request.Body)
for i := 1; ; i++ {
n, err := body.Read(buf)
if n > 0 {
select {
case streamCh <- string(buf[:n]):
case <-done:
return
}
}
if err != nil {
if err == io.EOF {
log.ZDebug(c, "read request body stream msg done", "clientMsgID", clientMsgID)
} else {
log.ZError(c, "read request body stream msg failed", err, "clientMsgID", clientMsgID, "error", err)
}
return
}
if n < 10 {
time.Sleep(time.Millisecond * 10)
}
}
}()

var (
packet []string
end bool
index int
errCount int
lastErr error
)
defer func() {
close(done)
if lastErr == nil {
apiresp.GinSuccess(c, nil)
} else {
m.putErr(c, lastErr)
}
}()
doAppend := func() {
if end == false && len(packet) == 0 {
return
}
ctx, cancel := context.WithTimeout(c, time.Second*10)
defer cancel()
req := &msg.AppendStreamMsgReq{
ConversationID: conversationID,
ClientMsgID: clientMsgID,
StartIndex: int64(index),
Packets: packet,
End: end,
}
_, lastErr = m.Client.AppendStreamMsg(ctx, req)
if lastErr == nil {
log.ZDebug(ctx, "AppendStreamMsg ok", "clientMsgID", clientMsgID)
index += len(packet)
packet = packet[:0]
errCount = 0
return
}
errCount++
if errs.ErrRecordNotFound.Is(lastErr) {
log.ZWarn(c, "msg not found", nil, "clientMsgID", clientMsgID)
return
} else if errs.ErrNoPermission.Is(lastErr) {
log.ZError(c, "msg permission error", nil, "clientMsgID", clientMsgID)
return
} else {
log.ZError(c, "append stream msg failed", lastErr, "clientMsgID", clientMsgID, "errCount", errCount)
time.Sleep(time.Millisecond * 50 * time.Duration(errCount))
}
}
for errCount < 10 {
select {
case s, ok := <-streamCh:
if ok {
packet = append(packet, s)
}
if !ok {
end = true
}
doAppend()
if end == true && lastErr == nil {
return
}
}
}
}

func NewUTF8Reader(r io.Reader) io.Reader {
return &UTF8Reader{
r: bufio.NewReaderSize(r, 512),
}
}

type UTF8Reader struct {
r *bufio.Reader
buf bytes.Buffer
}

func (r *UTF8Reader) Read(b []byte) (int, error) {
for {
n, err := r.r.Read(b)
if err != nil {
return 0, err
}
r.buf.Write(b[:n])
data := r.buf.Bytes()
minIndex := min(len(b), len(data))
if minIndex == 0 {
continue
}
for i := minIndex; i > 0; i-- {
if utf8.Valid(data[:i]) {
n, err := r.buf.Read(b[:i])
if err != nil {
return 0, err
}
if n != i {
return 0, fmt.Errorf("invalid UTF-8 encoding")
}
return n, nil
}
}
}
}

func (m *MessageApi) putErr(c *gin.Context, err error) {
c.JSON(http.StatusOK, apiresp.ParseError(err))
}
Loading
Loading