Files
meowlib/client/helpers/messageHelper.go
ycc 66a6674a6a
Some checks failed
continuous-integration/drone/push Build is failing
ProcessSentMessages added for server acks
2026-02-28 10:08:55 +01:00

173 lines
5.9 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package helpers
import (
"os"
"path/filepath"
"strconv"
"strings"
"forge.redroom.link/yves/meowlib"
"forge.redroom.link/yves/meowlib/client"
)
func PackMessageForServer(packedMsg *meowlib.PackedUserMessage, srvuid string) ([]byte, string, error) {
// Get the message server
srv, err := client.GetConfig().GetIdentity().MessageServers.LoadServer(srvuid)
if err != nil {
return nil, "messageBuildPostprocess : LoadServer", err
}
// Creating Server message for transporting the user message
toServerMessage := srv.BuildToServerMessageFromUserMessage(packedMsg)
data, err := srv.ProcessOutboundMessage(toServerMessage)
if err != nil {
return nil, "messageBuildPostprocess : ProcessOutboundMessage", err
}
return data, "", nil
}
func CreateStorePackUserMessageForServer(message string, srvuid string, peer_uid string, replyToUid string, filelist []string) ([]byte, string, error) {
usermessage, errtxt, err := CreateAndStoreUserMessage(message, peer_uid, replyToUid, filelist)
if err != nil {
return nil, errtxt, err
}
return PackMessageForServer(usermessage, srvuid)
}
func CreateAndStoreUserMessage(message string, peer_uid string, replyToUid string, filelist []string) (*meowlib.PackedUserMessage, string, error) {
peer := client.GetConfig().GetIdentity().Peers.GetFromUid(peer_uid)
// Creating User message
usermessage, err := peer.BuildSimpleUserMessage([]byte(message))
if err != nil {
return nil, "PrepareServerMessage : BuildSimpleUserMessage", err
}
for _, file := range filelist {
err = usermessage.AddFile(file, client.GetConfig().Chunksize)
if err != nil {
return nil, "PrepareServerMessage : AddFile", err
}
}
usermessage.Status.AnswerToUuid = replyToUid
// Store message
err = peer.StoreMessage(usermessage, nil)
if err != nil {
return nil, "messageBuildPostprocess : StoreMessage", err
}
// Prepare cyphered + packed user message
packedMsg, err := peer.ProcessOutboundUserMessage(usermessage)
if err != nil {
return nil, "messageBuildPostprocess : ProcessOutboundUserMessage", err
}
return packedMsg, "", nil
}
func BuildAckMessage(messageUid string, srvuid string, peer_uid string, received int64, processed int64) ([]byte, string, error) {
peer := client.GetConfig().GetIdentity().Peers.GetFromUid(peer_uid)
srv, err := client.GetConfig().GetIdentity().MessageServers.LoadServer(srvuid)
if err != nil {
return nil, "PrepareServerMessage : LoadServer", err
}
// Creating User message
usermessage, err := peer.BuildSimpleUserMessage(nil)
if err != nil {
return nil, "PrepareServerMessage : BuildSimpleUserMessage", err
}
usermessage.Status.Uuid = messageUid
usermessage.Status.Received = uint64(received)
usermessage.Status.Processed = uint64(processed)
// Prepare cyphered + packed user message
packedMsg, err := peer.ProcessOutboundUserMessage(usermessage)
if err != nil {
return nil, "PrepareServerMessage : ProcessOutboundUserMessage", err
}
// Creating Server message for transporting the user message
toServerMessage := srv.BuildToServerMessageFromUserMessage(packedMsg)
data, err := srv.ProcessOutboundMessage(toServerMessage)
if err != nil {
return nil, "PrepareServerMessage : ProcessOutboundMessage", err
}
return data, "", nil
}
func ReadAckMessageResponse() {
//! update the status in message store
}
// ProcessSentMessages scans every send queue under storagePath/queues/, updates
// the message storage entry with server delivery info for each sent job, then
// removes the job from the queue. Returns the number of messages updated.
//
// Callers must follow two conventions when building a SendJob:
// - job.Queue = peer UID (used to look up the peer and its DB files)
// - job.File = a path whose basename without extension is the message's
// SQLite row ID as a decimal integer (e.g. "42.bin")
func ProcessSentMessages(storagePath, password string) int {
queueDir := filepath.Join(storagePath, "queues")
entries, err := os.ReadDir(queueDir)
if err != nil {
logger.Warn().Err(err).Str("dir", queueDir).Msg("ProcessSentMessages: ReadDir")
return 0
}
updated := 0
identity := client.GetConfig().GetIdentity()
for _, entry := range entries {
if entry.IsDir() {
continue
}
queue := entry.Name()
jobs, err := client.GetSentJobs(storagePath, queue)
if err != nil {
logger.Error().Err(err).Str("queue", queue).Msg("ProcessSentMessages: GetSentJobs")
continue
}
for _, job := range jobs {
if job.SuccessfulServer == nil || job.SentAt == nil {
// No delivery info discard the job so it doesn't block the queue
if err := client.DeleteSendJob(storagePath, queue, job.ID); err != nil {
logger.Error().Err(err).Int64("id", job.ID).Msg("ProcessSentMessages: DeleteSendJob (incomplete)")
}
continue
}
// Resolve the peer from the queue name to get its DB file list
peer := identity.Peers.GetFromUid(queue)
if peer == nil || len(peer.DbIds) == 0 {
logger.Warn().Str("queue", queue).Msg("ProcessSentMessages: peer not found or has no DB")
continue
}
dbFile := peer.DbIds[len(peer.DbIds)-1]
// Parse the DB row ID from the job file's basename (e.g. "42.bin" → 42)
base := strings.TrimSuffix(filepath.Base(job.File), filepath.Ext(job.File))
dbId, err := strconv.ParseInt(base, 10, 64)
if err != nil {
logger.Error().Err(err).Str("file", job.File).Msg("ProcessSentMessages: parse dbId from filename")
continue
}
serverUid := job.Servers[*job.SuccessfulServer].GetUid()
receiveTime := uint64(job.SentAt.Unix())
if err := client.SetMessageServerDelivery(dbFile, dbId, serverUid, receiveTime, password); err != nil {
logger.Error().Err(err).Str("queue", queue).Int64("dbId", dbId).Msg("ProcessSentMessages: SetMessageServerDelivery")
continue
}
if err := client.DeleteSendJob(storagePath, queue, job.ID); err != nil {
logger.Error().Err(err).Int64("id", job.ID).Msg("ProcessSentMessages: DeleteSendJob")
}
updated++
}
}
return updated
}