Merge pull request 'v1.2.1' (#5) from v1.2.1 into main
Reviewed-on: #5
This commit was merged in pull request #5.
This commit is contained in:
+2
-2
@@ -2,7 +2,7 @@
|
|||||||
|
|
||||||
julia_version = "1.12.6"
|
julia_version = "1.12.6"
|
||||||
manifest_format = "2.0"
|
manifest_format = "2.0"
|
||||||
project_hash = "866f6d0804412d52eacd6423616500484f0060f0"
|
project_hash = "077e6beb148583a11d2d65eed055b8ff4a7db4b6"
|
||||||
|
|
||||||
[[deps.AliasTables]]
|
[[deps.AliasTables]]
|
||||||
deps = ["PtrArrays", "Random"]
|
deps = ["PtrArrays", "Random"]
|
||||||
@@ -812,7 +812,7 @@ version = "1.0.21+0"
|
|||||||
deps = ["Arrow", "Base64", "DataFrames", "Dates", "GeneralUtils", "HTTP", "JSON", "NATS", "PrettyPrinting", "Revise", "UUIDs"]
|
deps = ["Arrow", "Base64", "DataFrames", "Dates", "GeneralUtils", "HTTP", "JSON", "NATS", "PrettyPrinting", "Revise", "UUIDs"]
|
||||||
path = "."
|
path = "."
|
||||||
uuid = "f2724d33-f338-4a57-b9f8-1be882570d10"
|
uuid = "f2724d33-f338-4a57-b9f8-1be882570d10"
|
||||||
version = "0.5.6"
|
version = "1.2.1"
|
||||||
|
|
||||||
[[deps.nghttp2_jll]]
|
[[deps.nghttp2_jll]]
|
||||||
deps = ["Artifacts", "Libdl"]
|
deps = ["Artifacts", "Libdl"]
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
name = "msghandler"
|
name = "msghandler"
|
||||||
uuid = "f2724d33-f338-4a57-b9f8-1be882570d10"
|
uuid = "f2724d33-f338-4a57-b9f8-1be882570d10"
|
||||||
version = "1.2.0"
|
version = "1.2.1"
|
||||||
authors = ["narawat <narawat@gmail.com>"]
|
authors = ["narawat <narawat@gmail.com>"]
|
||||||
|
|
||||||
[deps]
|
[deps]
|
||||||
|
|||||||
+24
-13
@@ -452,19 +452,20 @@ function smartpack(
|
|||||||
payloads = msg_payload_v1[]
|
payloads = msg_payload_v1[]
|
||||||
for (dataname, payload_data, payload_type) in data
|
for (dataname, payload_data, payload_type) in data
|
||||||
# @show dataname typeof(payload_data)
|
# @show dataname typeof(payload_data)
|
||||||
|
# @info "msghandler smartpack() 1" @__LINE__
|
||||||
|
|
||||||
# Serialize data based on type. use bytes as medium for every datatype
|
# Serialize data based on type. use bytes as medium for every datatype
|
||||||
payload_bytes = _serialize_data(payload_data, payload_type)
|
payload_bytes = _serialize_data(payload_data, payload_type)
|
||||||
|
# @info "msghandler smartpack() 2" @__LINE__
|
||||||
payload_size = length(payload_bytes) # Calculate payload size in bytes
|
payload_size = length(payload_bytes) # Calculate payload size in bytes
|
||||||
log_trace(correlation_id, "Serialized payload '$dataname' (payload_type: $payload_type) size: $payload_size bytes") # Log payload size
|
# log_trace(correlation_id, "Serialized payload '$dataname' (payload_type: $payload_type) size: $payload_size bytes") # Log payload size
|
||||||
|
|
||||||
# Decision: Direct vs Link
|
# Decision: Direct vs Link
|
||||||
if payload_size < size_threshold # Check if payload is small enough for direct transport
|
if payload_size < size_threshold # Check if payload is small enough for direct transport
|
||||||
# Direct path - Base64 encode and include in message envelope
|
# Direct path - Base64 encode and include in message envelope
|
||||||
payload_b64 = Base64.base64encode(payload_bytes) # Encode bytes as base64 string
|
payload_b64 = Base64.base64encode(payload_bytes) # Encode bytes as base64 string
|
||||||
log_trace(correlation_id, "Using direct transport for $payload_size bytes") # Log transport choice
|
# log_trace(correlation_id, "Using direct transport for $payload_size bytes") # Log transport choice
|
||||||
|
# @info "msghandler smartpack() 3" @__LINE__
|
||||||
# Determine encoding based on payload_type
|
# Determine encoding based on payload_type
|
||||||
encoding = "base64"
|
encoding = "base64"
|
||||||
if payload_type == "jsontable"
|
if payload_type == "jsontable"
|
||||||
@@ -486,8 +487,9 @@ function smartpack(
|
|||||||
)
|
)
|
||||||
push!(payloads, payload)
|
push!(payloads, payload)
|
||||||
else
|
else
|
||||||
|
# @info "msghandler smartpack() 4" @__LINE__
|
||||||
# Link path - Upload to HTTP server, include URL in message envelope
|
# Link path - Upload to HTTP server, include URL in message envelope
|
||||||
log_trace(correlation_id, "Using link transport, uploading to fileserver") # Log link transport choice
|
# log_trace(correlation_id, "Using link transport, uploading to fileserver") # Log link transport choice
|
||||||
|
|
||||||
# Upload to HTTP server
|
# Upload to HTTP server
|
||||||
response = fileserver_upload_handler(fileserver_url, dataname, payload_bytes)
|
response = fileserver_upload_handler(fileserver_url, dataname, payload_bytes)
|
||||||
@@ -497,7 +499,7 @@ function smartpack(
|
|||||||
end
|
end
|
||||||
|
|
||||||
url = response["url"] # URL for the uploaded data
|
url = response["url"] # URL for the uploaded data
|
||||||
log_trace(correlation_id, "Uploaded to URL: $url") # Log successful upload
|
# log_trace(correlation_id, "Uploaded to URL: $url") # Log successful upload
|
||||||
|
|
||||||
# Determine encoding based on payload_type
|
# Determine encoding based on payload_type
|
||||||
encoding = "none"
|
encoding = "none"
|
||||||
@@ -520,6 +522,7 @@ function smartpack(
|
|||||||
)
|
)
|
||||||
push!(payloads, payload)
|
push!(payloads, payload)
|
||||||
end
|
end
|
||||||
|
# @info "msghandler smartpack() 5" @__LINE__
|
||||||
end
|
end
|
||||||
|
|
||||||
# Create msg_envelope_v1 with all payloads
|
# Create msg_envelope_v1 with all payloads
|
||||||
@@ -549,7 +552,7 @@ function smartpack(
|
|||||||
# # Publish message to NATS using existing connection
|
# # Publish message to NATS using existing connection
|
||||||
# publish_message(NATS_connection, subject, env_json_str, correlation_id)
|
# publish_message(NATS_connection, subject, env_json_str, correlation_id)
|
||||||
# end
|
# end
|
||||||
|
# @info "msghandler smartpack() 6" @__LINE__
|
||||||
return (env, env_json_str)
|
return (env, env_json_str)
|
||||||
end
|
end
|
||||||
|
|
||||||
@@ -637,21 +640,22 @@ function _serialize_data(data::Any, payload_type::String)
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
if payload_type == "text" # Text data - convert to UTF-8 bytes
|
if payload_type == "text" # Text data - convert to UTF-8 bytes
|
||||||
if isa(data, String)
|
# @info "msghandler _serialize_data() 1" @__LINE__
|
||||||
data_bytes = Vector{UInt8}(data) # Convert string to UTF-8 bytes
|
data_bytes = Vector{UInt8}(string(data)) # Convert string to UTF-8 bytes
|
||||||
return data_bytes
|
return data_bytes
|
||||||
else
|
|
||||||
error("Text data must be a String")
|
|
||||||
end
|
|
||||||
elseif payload_type == "dictionary" # JSON data - serialize directly
|
elseif payload_type == "dictionary" # JSON data - serialize directly
|
||||||
|
# @info "msghandler _serialize_data() 2" @__LINE__
|
||||||
json_str = JSON.json(data) # Convert Julia data to JSON string
|
json_str = JSON.json(data) # Convert Julia data to JSON string
|
||||||
json_str_bytes = Vector{UInt8}(json_str) # Convert JSON string to bytes
|
json_str_bytes = Vector{UInt8}(json_str) # Convert JSON string to bytes
|
||||||
return json_str_bytes
|
return json_str_bytes
|
||||||
elseif payload_type == "arrowtable" # Arrow table data - convert to Arrow IPC stream
|
elseif payload_type == "arrowtable" # Arrow table data - convert to Arrow IPC stream
|
||||||
|
# @info "msghandler _serialize_data() 3" @__LINE__
|
||||||
io = IOBuffer() # Create in-memory buffer
|
io = IOBuffer() # Create in-memory buffer
|
||||||
Arrow.write(io, data) # Write data as Arrow IPC stream to buffer
|
Arrow.write(io, data) # Write data as Arrow IPC stream to buffer
|
||||||
return take!(io) # Return the buffer contents as bytes
|
return take!(io) # Return the buffer contents as bytes
|
||||||
elseif payload_type == "jsontable" # JSON table data - convert to JSON
|
elseif payload_type == "jsontable" # JSON table data - convert to JSON
|
||||||
|
# @info "msghandler _serialize_data() 4" @__LINE__
|
||||||
|
|
||||||
# data can be Vector{NamedTuple}, Vector{Dict}, or DataFrame
|
# data can be Vector{NamedTuple}, Vector{Dict}, or DataFrame
|
||||||
# If DataFrame, convert to Vector{Dict} first
|
# If DataFrame, convert to Vector{Dict} first
|
||||||
if isa(data, DataFrame)
|
if isa(data, DataFrame)
|
||||||
@@ -667,29 +671,35 @@ function _serialize_data(data::Any, payload_type::String)
|
|||||||
json_str = JSON.json(rows)
|
json_str = JSON.json(rows)
|
||||||
return Vector{UInt8}(json_str)
|
return Vector{UInt8}(json_str)
|
||||||
else
|
else
|
||||||
|
# @info "msghandler _serialize_data() 5" @__LINE__
|
||||||
|
|
||||||
# Already Vector{NamedTuple} or Vector{Dict}
|
# Already Vector{NamedTuple} or Vector{Dict}
|
||||||
json_str = JSON.json(data)
|
json_str = JSON.json(data)
|
||||||
return Vector{UInt8}(json_str)
|
return Vector{UInt8}(json_str)
|
||||||
end
|
end
|
||||||
elseif payload_type == "image" # Image data - treat as binary
|
elseif payload_type == "image" # Image data - treat as binary
|
||||||
|
# @info "msghandler _serialize_data() 6" @__LINE__
|
||||||
if isa(data, Vector{UInt8})
|
if isa(data, Vector{UInt8})
|
||||||
return data # Return binary data directly
|
return data # Return binary data directly
|
||||||
else
|
else
|
||||||
error("Image data must be Vector{UInt8}")
|
error("Image data must be Vector{UInt8}")
|
||||||
end
|
end
|
||||||
elseif payload_type == "audio" # Audio data - treat as binary
|
elseif payload_type == "audio" # Audio data - treat as binary
|
||||||
|
# @info "msghandler _serialize_data() 7" @__LINE__
|
||||||
if isa(data, Vector{UInt8})
|
if isa(data, Vector{UInt8})
|
||||||
return data # Return binary data directly
|
return data # Return binary data directly
|
||||||
else
|
else
|
||||||
error("Audio data must be Vector{UInt8}")
|
error("Audio data must be Vector{UInt8}")
|
||||||
end
|
end
|
||||||
elseif payload_type == "video" # Video data - treat as binary
|
elseif payload_type == "video" # Video data - treat as binary
|
||||||
|
# @info "msghandler _serialize_data() 8" @__LINE__
|
||||||
if isa(data, Vector{UInt8})
|
if isa(data, Vector{UInt8})
|
||||||
return data # Return binary data directly
|
return data # Return binary data directly
|
||||||
else
|
else
|
||||||
error("Video data must be Vector{UInt8}")
|
error("Video data must be Vector{UInt8}")
|
||||||
end
|
end
|
||||||
elseif payload_type == "binary" # Binary data - treat as binary
|
elseif payload_type == "binary" # Binary data - treat as binary
|
||||||
|
# @info "msghandler _serialize_data() 9" @__LINE__
|
||||||
if isa(data, IOBuffer) # Check if data is an IOBuffer
|
if isa(data, IOBuffer) # Check if data is an IOBuffer
|
||||||
return take!(data) # Return buffer contents as bytes
|
return take!(data) # Return buffer contents as bytes
|
||||||
elseif isa(data, Vector{UInt8}) # Check if data is already binary
|
elseif isa(data, Vector{UInt8}) # Check if data is already binary
|
||||||
@@ -698,6 +708,7 @@ function _serialize_data(data::Any, payload_type::String)
|
|||||||
error("Binary data must be binary (Vector{UInt8} or IOBuffer)")
|
error("Binary data must be binary (Vector{UInt8} or IOBuffer)")
|
||||||
end
|
end
|
||||||
else # Unknown type
|
else # Unknown type
|
||||||
|
# @info "msghandler _serialize_data() 10" @__LINE__
|
||||||
error("Unknown payload_type: $payload_type")
|
error("Unknown payload_type: $payload_type")
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|||||||
Reference in New Issue
Block a user