2024-09-20 17:43:56 +02:00
|
|
|
import
|
|
|
|
std/[times, strutils, asyncnet, os, sequtils],
|
|
|
|
results,
|
|
|
|
chronos,
|
|
|
|
metrics,
|
|
|
|
re,
|
|
|
|
chronicles
|
2024-09-20 13:23:53 +02:00
|
|
|
import ./query_metrics
|
2023-06-07 10:08:43 +02:00
|
|
|
|
2024-07-09 13:14:28 +02:00
|
|
|
include db_connector/db_postgres
|
2023-06-07 10:08:43 +02:00
|
|
|
|
2023-12-14 07:16:39 +01:00
|
|
|
type DataProc* = proc(result: ptr PGresult) {.closure, gcsafe, raises: [].}
|
2023-10-31 14:46:46 +01:00
|
|
|
|
2023-06-07 10:08:43 +02:00
|
|
|
## Connection management
|
|
|
|
|
|
|
|
proc check*(db: DbConn): Result[void, string] =
|
|
|
|
var message: string
|
|
|
|
try:
|
|
|
|
message = $db.pqErrorMessage()
|
2024-03-16 00:08:47 +01:00
|
|
|
except ValueError, DbError:
|
2023-06-07 10:08:43 +02:00
|
|
|
return err("exception in check: " & getCurrentExceptionMsg())
|
|
|
|
|
|
|
|
if message.len > 0:
|
|
|
|
return err($message)
|
|
|
|
|
|
|
|
return ok()
|
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
proc open*(connString: string): Result[DbConn, string] =
|
2023-06-07 10:08:43 +02:00
|
|
|
## Opens a new connection.
|
|
|
|
var conn: DbConn = nil
|
|
|
|
try:
|
2024-03-16 00:08:47 +01:00
|
|
|
conn = open("", "", "", connString)
|
2023-06-07 10:08:43 +02:00
|
|
|
except DbError:
|
2024-03-16 00:08:47 +01:00
|
|
|
return err("exception opening new connection: " & getCurrentExceptionMsg())
|
2023-06-07 10:08:43 +02:00
|
|
|
|
|
|
|
if conn.status != CONNECTION_OK:
|
|
|
|
let checkRes = conn.check()
|
|
|
|
if checkRes.isErr():
|
|
|
|
return err("failed to connect to database: " & checkRes.error)
|
|
|
|
|
|
|
|
return err("unknown reason")
|
|
|
|
|
2024-09-06 11:33:15 +02:00
|
|
|
## registering the socket fd in chronos for better wait for data
|
|
|
|
let asyncFd = cast[asyncengine.AsyncFD](pqsocket(conn))
|
|
|
|
asyncengine.register(asyncFd)
|
|
|
|
|
|
|
|
return ok(conn)
|
|
|
|
|
|
|
|
proc closeDbConn*(db: DbConn) {.raises: [OSError].} =
|
|
|
|
let fd = db.pqsocket()
|
|
|
|
if fd != -1:
|
|
|
|
asyncengine.unregister(cast[asyncengine.AsyncFD](fd))
|
|
|
|
db.close()
|
2023-06-07 10:08:43 +02:00
|
|
|
|
2024-09-20 13:23:53 +02:00
|
|
|
proc `$`(self: SqlQuery): string =
|
|
|
|
return cast[string](self)
|
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
proc sendQuery(
|
|
|
|
db: DbConn, query: SqlQuery, args: seq[string]
|
|
|
|
): Future[Result[void, string]] {.async.} =
|
2023-10-31 14:46:46 +01:00
|
|
|
## This proc can be used directly for queries that don't retrieve values back.
|
2023-06-07 10:08:43 +02:00
|
|
|
|
|
|
|
if db.status != CONNECTION_OK:
|
2023-11-07 13:38:37 +01:00
|
|
|
db.check().isOkOr:
|
|
|
|
return err("failed to connect to database: " & $error)
|
2023-06-07 10:08:43 +02:00
|
|
|
|
|
|
|
return err("unknown reason")
|
|
|
|
|
|
|
|
var wellFormedQuery = ""
|
|
|
|
try:
|
|
|
|
wellFormedQuery = dbFormat(query, args)
|
|
|
|
except DbError:
|
2024-03-16 00:08:47 +01:00
|
|
|
return err("exception formatting the query: " & getCurrentExceptionMsg())
|
2023-06-07 10:08:43 +02:00
|
|
|
|
|
|
|
let success = db.pqsendQuery(cstring(wellFormedQuery))
|
|
|
|
if success != 1:
|
2023-11-07 13:38:37 +01:00
|
|
|
db.check().isOkOr:
|
|
|
|
return err("failed pqsendQuery: " & $error)
|
2023-06-07 10:08:43 +02:00
|
|
|
|
|
|
|
return err("failed pqsendQuery: unknown reason")
|
|
|
|
|
2023-10-31 14:46:46 +01:00
|
|
|
return ok()
|
|
|
|
|
2023-11-07 13:38:37 +01:00
|
|
|
proc sendQueryPrepared(
|
2024-03-16 00:08:47 +01:00
|
|
|
db: DbConn,
|
|
|
|
stmtName: string,
|
|
|
|
paramValues: openArray[string],
|
|
|
|
paramLengths: openArray[int32],
|
|
|
|
paramFormats: openArray[int32],
|
|
|
|
): Result[void, string] {.raises: [].} =
|
2023-11-07 13:38:37 +01:00
|
|
|
## This proc can be used directly for queries that don't retrieve values back.
|
|
|
|
|
|
|
|
if paramValues.len != paramLengths.len or paramValues.len != paramFormats.len or
|
2024-03-16 00:08:47 +01:00
|
|
|
paramLengths.len != paramFormats.len:
|
|
|
|
let lengthsErrMsg =
|
|
|
|
$paramValues.len & " " & $paramLengths.len & " " & $paramFormats.len
|
2023-11-07 13:38:37 +01:00
|
|
|
return err("lengths discrepancies in sendQueryPrepared: " & $lengthsErrMsg)
|
|
|
|
|
|
|
|
if db.status != CONNECTION_OK:
|
|
|
|
db.check().isOkOr:
|
|
|
|
return err("failed to connect to database: " & $error)
|
|
|
|
|
|
|
|
return err("unknown reason")
|
|
|
|
|
|
|
|
var cstrArrayParams = allocCStringArray(paramValues)
|
2024-03-16 00:08:47 +01:00
|
|
|
defer:
|
|
|
|
deallocCStringArray(cstrArrayParams)
|
2023-11-07 13:38:37 +01:00
|
|
|
|
|
|
|
let nParams = cast[int32](paramValues.len)
|
|
|
|
|
|
|
|
const ResultFormat = 0 ## 0 for text format, 1 for binary format.
|
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
let success = db.pqsendQueryPrepared(
|
|
|
|
stmtName,
|
|
|
|
nParams,
|
|
|
|
cstrArrayParams,
|
|
|
|
unsafeAddr paramLengths[0],
|
|
|
|
unsafeAddr paramFormats[0],
|
|
|
|
ResultFormat,
|
|
|
|
)
|
2023-11-07 13:38:37 +01:00
|
|
|
if success != 1:
|
|
|
|
db.check().isOkOr:
|
|
|
|
return err("failed pqsendQueryPrepared: " & $error)
|
|
|
|
|
|
|
|
return err("failed pqsendQueryPrepared: unknown reason")
|
|
|
|
|
|
|
|
return ok()
|
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
proc waitQueryToFinish(
|
|
|
|
db: DbConn, rowCallback: DataProc = nil
|
|
|
|
): Future[Result[void, string]] {.async.} =
|
2023-10-31 14:46:46 +01:00
|
|
|
## The 'rowCallback' param is != nil when the underlying query wants to retrieve results (SELECT.)
|
|
|
|
## For other queries, like "INSERT", 'rowCallback' should be nil.
|
2023-06-07 10:08:43 +02:00
|
|
|
|
2024-09-06 11:33:15 +02:00
|
|
|
var dataAvailable = false
|
|
|
|
proc onDataAvailable(udata: pointer) {.gcsafe, raises: [].} =
|
|
|
|
dataAvailable = true
|
|
|
|
|
|
|
|
let asyncFd = cast[asyncengine.AsyncFD](pqsocket(db))
|
2023-06-07 10:08:43 +02:00
|
|
|
|
2024-09-06 11:33:15 +02:00
|
|
|
asyncengine.addReader2(asyncFd, onDataAvailable).isOkOr:
|
|
|
|
return err("failed to add event reader in waitQueryToFinish: " & $error)
|
2023-06-07 10:08:43 +02:00
|
|
|
|
2024-09-06 11:33:15 +02:00
|
|
|
while not dataAvailable:
|
|
|
|
await sleepAsync(timer.milliseconds(1))
|
2023-06-07 10:08:43 +02:00
|
|
|
|
2023-11-07 13:38:37 +01:00
|
|
|
## Now retrieve the result
|
|
|
|
while true:
|
2023-10-31 14:46:46 +01:00
|
|
|
let pqResult = db.pqgetResult()
|
2023-11-07 13:38:37 +01:00
|
|
|
|
2023-06-07 10:08:43 +02:00
|
|
|
if pqResult == nil:
|
2023-11-07 13:38:37 +01:00
|
|
|
db.check().isOkOr:
|
|
|
|
return err("error in query: " & $error)
|
2023-06-07 10:08:43 +02:00
|
|
|
|
2023-10-31 14:46:46 +01:00
|
|
|
return ok() # reached the end of the results
|
2023-06-07 10:08:43 +02:00
|
|
|
|
2023-10-31 14:46:46 +01:00
|
|
|
if not rowCallback.isNil():
|
|
|
|
rowCallback(pqResult)
|
2023-06-07 10:08:43 +02:00
|
|
|
|
|
|
|
pqclear(pqResult)
|
2023-10-31 14:46:46 +01:00
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
proc dbConnQuery*(
|
|
|
|
db: DbConn, query: SqlQuery, args: seq[string], rowCallback: DataProc
|
|
|
|
): Future[Result[void, string]] {.async, gcsafe.} =
|
2024-09-20 13:23:53 +02:00
|
|
|
let cleanedQuery = ($query).replace(" ", "").replace("\n", "")
|
|
|
|
## remove everything between ' or " all possible sequence of numbers. e.g. rm partition partition
|
|
|
|
var querySummary = cleanedQuery.replace(re"""(['"]).*?\1""", "")
|
|
|
|
querySummary = querySummary.replace(re"\d+", "")
|
|
|
|
querySummary = "query_tag_" & querySummary[0 ..< min(querySummary.len, 200)]
|
|
|
|
|
|
|
|
var queryStartTime = getTime().toUnixFloat()
|
|
|
|
|
2023-10-31 14:46:46 +01:00
|
|
|
(await db.sendQuery(query, args)).isOkOr:
|
|
|
|
return err("error in dbConnQuery calling sendQuery: " & $error)
|
|
|
|
|
2024-09-20 17:43:56 +02:00
|
|
|
let sendDuration = getTime().toUnixFloat() - queryStartTime
|
|
|
|
query_time_secs.set(sendDuration, [querySummary, "sendQuery"])
|
2024-09-20 13:23:53 +02:00
|
|
|
|
|
|
|
queryStartTime = getTime().toUnixFloat()
|
|
|
|
|
2023-10-31 14:46:46 +01:00
|
|
|
(await db.waitQueryToFinish(rowCallback)).isOkOr:
|
|
|
|
return err("error in dbConnQuery calling waitQueryToFinish: " & $error)
|
|
|
|
|
2024-09-20 17:43:56 +02:00
|
|
|
let waitDuration = getTime().toUnixFloat() - queryStartTime
|
|
|
|
query_time_secs.set(waitDuration, [querySummary, "waitFinish"])
|
2024-09-20 13:23:53 +02:00
|
|
|
|
|
|
|
query_count.inc(labelValues = [querySummary])
|
|
|
|
|
2024-09-20 17:43:56 +02:00
|
|
|
if "insert" notin ($query).toLower():
|
|
|
|
debug "dbConnQuery",
|
|
|
|
query = $query,
|
|
|
|
querySummary,
|
|
|
|
waitDurationSecs = waitDuration,
|
|
|
|
sendDurationSecs = sendDuration
|
|
|
|
|
2023-10-31 14:46:46 +01:00
|
|
|
return ok()
|
2023-11-07 13:38:37 +01:00
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
proc dbConnQueryPrepared*(
|
|
|
|
db: DbConn,
|
|
|
|
stmtName: string,
|
|
|
|
paramValues: seq[string],
|
|
|
|
paramLengths: seq[int32],
|
|
|
|
paramFormats: seq[int32],
|
|
|
|
rowCallback: DataProc,
|
|
|
|
): Future[Result[void, string]] {.async, gcsafe.} =
|
2024-09-20 13:23:53 +02:00
|
|
|
var queryStartTime = getTime().toUnixFloat()
|
2024-03-16 00:08:47 +01:00
|
|
|
db.sendQueryPrepared(stmtName, paramValues, paramLengths, paramFormats).isOkOr:
|
2023-11-07 13:38:37 +01:00
|
|
|
return err("error in dbConnQueryPrepared calling sendQuery: " & $error)
|
|
|
|
|
2024-09-20 17:43:56 +02:00
|
|
|
let sendDuration = getTime().toUnixFloat() - queryStartTime
|
|
|
|
query_time_secs.set(sendDuration, [stmtName, "sendQuery"])
|
2024-09-20 13:23:53 +02:00
|
|
|
|
|
|
|
queryStartTime = getTime().toUnixFloat()
|
|
|
|
|
2023-11-07 13:38:37 +01:00
|
|
|
(await db.waitQueryToFinish(rowCallback)).isOkOr:
|
|
|
|
return err("error in dbConnQueryPrepared calling waitQueryToFinish: " & $error)
|
|
|
|
|
2024-09-20 17:43:56 +02:00
|
|
|
let waitDuration = getTime().toUnixFloat() - queryStartTime
|
|
|
|
query_time_secs.set(waitDuration, [stmtName, "waitFinish"])
|
2024-09-20 13:23:53 +02:00
|
|
|
|
|
|
|
query_count.inc(labelValues = [stmtName])
|
|
|
|
|
2024-09-20 17:43:56 +02:00
|
|
|
if "insert" notin stmtName.toLower():
|
|
|
|
debug "dbConnQueryPrepared",
|
|
|
|
stmtName, waitDurationSecs = waitDuration, sendDurationSecs = sendDuration
|
|
|
|
|
2023-11-07 13:38:37 +01:00
|
|
|
return ok()
|