diff --git a/README.md b/README.md index cc53f6c..fb27df6 100644 --- a/README.md +++ b/README.md @@ -4,7 +4,7 @@ **A multimodal database engine written in Nim — 100% native, zero dependencies.** -[![Version](https://img.shields.io/badge/version-1.1.7-blue.svg)](baradadb.nimble) +[![Version](https://img.shields.io/badge/version-1.1.8-blue.svg)](baradadb.nimble) [![Documentation](https://img.shields.io/badge/docs-2_languages-blue.svg)](docs/index.md) [![Stars](https://img.shields.io/github/stars/katehonz/barabaDB?style=social)](https://github.com/katehonz/barabaDB) @@ -1568,7 +1568,7 @@ reflects 100% completion across all major phases. ## Changelog -See [CHANGELOG.md](CHANGELOG.md) for full release history. The latest release (**v1.2.0**) introduces the Unified Search Engine with heap-optimized HNSW, segment-based inverted indexing, boolean queries, phrase/proximity search, n-gram fuzzy matching, faceted search, and Porter2 stemmers for 5 languages. +See [CHANGELOG.md](CHANGELOG.md) for full release history. Package version is **v1.1.8**. The **v1.2.0** line (Unreleased) adds core storage hardening (hash MemTable, WAL group commit, schema persistence, ARC/wire stability) and the Unified Search Engine (heap-optimized HNSW, segment inverted index, boolean/phrase/n-gram/facets, multi-language stemmers). ## License diff --git a/baradadb.nimble b/baradadb.nimble index d47012d..2be7758 100644 --- a/baradadb.nimble +++ b/baradadb.nimble @@ -2,7 +2,7 @@ version = "1.1.8" author = "BaraDB Team" description = "BaraDB — Multimodal database written in Nim" -license = "Apache-2.0" +license = "BSD-3-Clause" srcDir = "src" bin = @["baradadb", "baramcp"] binDir = "build" diff --git a/clients/nim/README.md b/clients/nim/README.md index 279073d..42beb2f 100644 --- a/clients/nim/README.md +++ b/clients/nim/README.md @@ -217,4 +217,4 @@ See `examples/ormin_basic.nim` for a full sample. ## License -Apache-2.0 +BSD-3-Clause diff --git a/clients/nim/baradb.nimble b/clients/nim/baradb.nimble index 4610ef2..5e14ca5 100644 --- a/clients/nim/baradb.nimble +++ b/clients/nim/baradb.nimble @@ -3,7 +3,7 @@ version = "1.2.0" author = "BaraDB Team" description = "Official Nim client for BaraDB — async binary protocol client" -license = "Apache-2.0" +license = "BSD-3-Clause" srcDir = "src" # Dependencies — only Nim stdlib, no server code diff --git a/clients/nim/src/baradb/client.nim b/clients/nim/src/baradb/client.nim index 9a4d8ca..802ff35 100644 --- a/clients/nim/src/baradb/client.nim +++ b/clients/nim/src/baradb/client.nim @@ -92,6 +92,26 @@ proc newClient*(config: ClientConfig = defaultConfig()): BaraClient = sendLock: initAsyncLock(), ) +# Aliases for older call sites / server test suite +proc defaultClientConfig*(): ClientConfig {.inline.} = defaultConfig() +proc newBaraClient*(config: ClientConfig = defaultConfig()): BaraClient {.inline.} = + newClient(config) + +proc parseConnectionString*(connStr: string): ClientConfig = + ## Parse space-separated key=value pairs (libpq-style subset). + result = defaultConfig() + for part in connStr.split(" "): + let kv = part.split("=", 1) + if kv.len == 2: + case kv[0].toLowerAscii() + of "host": result.host = kv[1] + of "port": result.port = parseInt(kv[1]) + of "database", "dbname": result.database = kv[1] + of "user", "username": result.username = kv[1] + of "password", "pass": result.password = kv[1] + of "connect_timeout", "timeout": result.timeoutMs = parseInt(kv[1]) + else: discard + proc nextId*(client: BaraClient): uint32 = inc client.requestId client.requestId diff --git a/src/barabadb/core/backup.nim b/src/barabadb/core/backup.nim index 6de7364..6d13d23 100644 --- a/src/barabadb/core/backup.nim +++ b/src/barabadb/core/backup.nim @@ -146,7 +146,7 @@ proc formatTimestamp*(ts: int64): string = try: let dt = fromUnix(ts) result = format(dt, "yyyy-MM-dd HH:mm:ss") - except: + except CatchableError: result = $ts proc parseBackupFilename*(filename: string): int64 = @@ -163,7 +163,7 @@ proc parseBackupFilename*(filename: string): int64 = result = 0 else: result = 0 - except: + except CatchableError: result = 0 proc getArchiveSize*(input: string): int64 = @@ -175,7 +175,7 @@ proc getArchiveSize*(input: string): int64 = if exitCode == 0: try: result = parseBiggestInt(strip(outStr)) - except: + except CatchableError: result = getFileSize(input) # fallback else: result = getFileSize(input) @@ -189,7 +189,7 @@ proc getFreeSpace*(path: string): int64 = if exitCode == 0: try: result = parseBiggestInt(strip(outStr)) - except: + except CatchableError: result = -1 else: result = -1 @@ -701,14 +701,14 @@ when isMainModule: of "input", "i": target = val of "keep", "k": try: keepCount = parseInt(val) - except: quit("ERROR: --keep must be a number", 1) + except CatchableError: quit("ERROR: --keep must be a number", 1) of "exclude", "e": excludes.add(val) of "level", "l": try: compression = parseInt(val) if compression < 0 or compression > 9: quit("ERROR: --level must be between 0 and 9", 1) - except: quit("ERROR: --level must be a number", 1) + except CatchableError: quit("ERROR: --level must be a number", 1) of "dry-run": dryRun = true of "force", "f": force = true of "online": online = true diff --git a/src/barabadb/core/gossip.nim b/src/barabadb/core/gossip.nim index daec162..7dd45e8 100644 --- a/src/barabadb/core/gossip.nim +++ b/src/barabadb/core/gossip.nim @@ -277,7 +277,7 @@ proc sendGossipUdp(gp: GossipProtocol, target: GossipNode, msg: GossipMessage) = let data = serialize(msg) sock.sendTo(target.host, Port(target.port), cast[string](data)) sock.close() - except: + except CatchableError: discard proc broadcastGossip(gp: GossipProtocol) = @@ -298,7 +298,7 @@ proc handleIncomingGossip(gp: GossipProtocol, data: string, senderAddr: string) let parts = host.split(":") host = parts[0] if parts[1].len > 0: - port = try: parseInt(parts[1]) except: gp.gossipPort + port = try: parseInt(parts[1]) except CatchableError: gp.gossipPort let newNode = GossipNode( id: msg.senderId, host: host, port: port, state: nsAlive, incarnation: msg.senderIncarnation, @@ -306,7 +306,7 @@ proc handleIncomingGossip(gp: GossipProtocol, data: string, senderAddr: string) ) gp.addMember(newNode) gp.applyGossipMessage(msg) - except: + except CatchableError: discard proc startHealthCheck*(gp: GossipProtocol, intervalMs: int = 1000) {.async.} = @@ -342,14 +342,14 @@ proc startGossipListener*(gp: GossipProtocol) {.async.} = # Recreate socket after too many errors try: gp.sock.close() - except: + except CatchableError: discard try: gp.sock = newAsyncSocket(AF_INET, SOCK_DGRAM, IPPROTO_UDP) gp.sock.setSockOpt(OptReuseAddr, true) gp.sock.bindAddr(Port(gp.gossipPort)) consecutiveErrors = 0 - except: + except CatchableError: break # Exponential backoff with cap let delayMs = min(baseRetryDelayMs * (1 shl min(consecutiveErrors, 6)), 5000) diff --git a/src/barabadb/core/httpserver.nim b/src/barabadb/core/httpserver.nim index 7658a53..8ca90fe 100644 --- a/src/barabadb/core/httpserver.nim +++ b/src/barabadb/core/httpserver.nim @@ -103,7 +103,7 @@ proc verifyToken*(server: HttpServer, tokenStr: string): (bool, string, string) let userId = token.claims["sub"].node.str let role = if "role" in token.claims: token.claims["role"].node.str else: "user" return (true, userId, role) - except: + except CatchableError: return (false, "", "") # ---------------------------------------------------------------------- diff --git a/src/barabadb/core/raft.nim b/src/barabadb/core/raft.nim index dd95479..ce0ab61 100644 --- a/src/barabadb/core/raft.nim +++ b/src/barabadb/core/raft.nim @@ -202,7 +202,7 @@ proc applyCommitted(node: RaftNode) = let parts = entry.command.split(":") if parts.len >= 3: let action = parts[1] - let txnId = try: parseUInt(parts[2]) except: 0'u64 + let txnId = try: parseUInt(parts[2]) except CatchableError: 0'u64 if action == "PREPARE" and node.onDistTxnPrepare != nil: discard node.onDistTxnPrepare(txnId, @[]) elif action == "COMMIT" and node.onDistTxnCommit != nil: @@ -564,7 +564,7 @@ proc connectToPeer(net: RaftNetwork, peerId: string) {.async.} = let sock = newAsyncSocket() await sock.connect(host, Port(port)) net.peerSockets[peerId] = sock - except: + except CatchableError: discard proc send*(net: RaftNetwork, peerId: string, msg: RaftMessage) {.async.} = @@ -577,7 +577,7 @@ proc send*(net: RaftNetwork, peerId: string, msg: RaftMessage) {.async.} = bigEndian32(addr header[0], unsafeAddr payloadLen) try: await net.peerSockets[peerId].send(cast[string](header) & cast[string](data)) - except: + except CatchableError: net.peerSockets.del(peerId) proc broadcast*(net: RaftNetwork, msgs: seq[RaftMessage]) {.async.} = @@ -615,9 +615,9 @@ proc receiveLoop(net: RaftNetwork, client: AsyncSocket) {.async.} = let msg = deserializeRaftMessage(payload) try: await net.processMessage(msg) - except: + except CatchableError: discard - except: + except CatchableError: discard finally: client.close() @@ -641,7 +641,7 @@ proc run*(net: RaftNetwork) {.async.} = try: let client = await net.socket.accept() asyncCheck net.receiveLoop(client) - except: + except CatchableError: break proc stop*(net: RaftNetwork) = diff --git a/src/barabadb/core/replication.nim b/src/barabadb/core/replication.nim index a76f96b..1d30702 100644 --- a/src/barabadb/core/replication.nim +++ b/src/barabadb/core/replication.nim @@ -266,12 +266,12 @@ proc healthCheck*(rm: ReplicationManager) = sock.readLine(response) if response.strip() != "PONG": connected = false - except: + except CatchableError: connected = false - except: + except CatchableError: connected = false finally: - try: sock.close() except: discard + try: sock.close() except CatchableError: discard if not connected: acquire(rm.lock) diff --git a/src/barabadb/core/server.nim b/src/barabadb/core/server.nim index 3192c55..d7dc24b 100644 --- a/src/barabadb/core/server.nim +++ b/src/barabadb/core/server.nim @@ -65,55 +65,53 @@ proc newServerWithRegistry*(config: BaraConfig, registry: DatabaseRegistry): Ser let tlsConfig = newTLSConfig(config.certFile, config.keyFile) tls = newTLSContext(tlsConfig) - # Initialize sharding - let shardRouter = newShardRouter() + # Initialize sharding / gossip. Server fields own the refs; locals used inside + # callback closures are {.cursor.} so ARC does not form uncollectable cycles + # (local + closure env + object callback fields). let localId = if config.raftNodeId.len > 0: config.raftNodeId else: "node-" & $config.port - let cm = newClusterMembership(shardRouter, localId) - - # Wire shard migration callbacks to LSM (use default database) - shardRouter.iterateKeys = proc(shardId: int): seq[(string, seq[byte])] {.gcsafe.} = - var entries: seq[(string, seq[byte])] = @[] - for (key, value) in db.scanAll(): - if shardRouter.getShard(key) == shardId: - entries.add((key, value)) - return entries - - shardRouter.storeKeys = proc(shardId: int, entries: seq[(string, seq[byte])]) {.gcsafe.} = - for (key, value) in entries: - db.put(key, value) - - shardRouter.deleteKeys = proc(keys: seq[string]) {.gcsafe.} = - for key in keys: - db.delete(key) - - # Initialize gossip let gossipPort = config.raftPort + 100 - let gp = newGossipProtocol(localId, config.address, config.port, gossipPort = gossipPort) - - # Wire gossip → cluster membership - gp.onJoin = proc(node: GossipNode) {.gcsafe.} = - cm.onNodeJoin(node.id, node.host, node.port) - - gp.onLeave = proc(nodeId: string) {.gcsafe.} = - cm.onNodeLeave(nodeId) - - gp.onSuspect = proc(nodeId: string) {.gcsafe.} = - cm.onNodeSuspect(nodeId) - - # Initialize rate limiter let rl = newRateLimiter(rlaTokenBucket, config.rateLimitGlobal, config.rateLimitPerClient) result = Server(config: config, running: false, db: db, ctx: ctx, registry: registry, txnManager: ctx.txnManager, distTxnManager: newDistTxnManager(), replicationManager: newReplicationManager(), - shardRouter: shardRouter, - clusterMembership: cm, - gossipProtocol: gp, + shardRouter: newShardRouter(), + clusterMembership: nil, + gossipProtocol: newGossipProtocol(localId, config.address, config.port, gossipPort = gossipPort), tls: tls, rateLimiter: rl) + result.clusterMembership = newClusterMembership(result.shardRouter, localId) initLock(result.activeConnectionsLock) + # Wire shard migration callbacks to LSM (default database) + block: + let shardRouter {.cursor.} = result.shardRouter + let dbRef {.cursor.} = db + shardRouter.iterateKeys = proc(shardId: int): seq[(string, seq[byte])] {.gcsafe.} = + var entries: seq[(string, seq[byte])] = @[] + for (key, value) in dbRef.scanAll(): + if shardRouter.getShard(key) == shardId: + entries.add((key, value)) + return entries + shardRouter.storeKeys = proc(shardId: int, entries: seq[(string, seq[byte])]) {.gcsafe.} = + for (key, value) in entries: + dbRef.put(key, value) + shardRouter.deleteKeys = proc(keys: seq[string]) {.gcsafe.} = + for key in keys: + dbRef.delete(key) + + # Wire gossip → cluster membership + block: + let gp {.cursor.} = result.gossipProtocol + let cm {.cursor.} = result.clusterMembership + gp.onJoin = proc(node: GossipNode) {.gcsafe.} = + cm.onNodeJoin(node.id, node.host, node.port) + gp.onLeave = proc(nodeId: string) {.gcsafe.} = + cm.onNodeLeave(nodeId) + gp.onSuspect = proc(nodeId: string) {.gcsafe.} = + cm.onNodeSuspect(nodeId) + proc newServerWithDb*(config: BaraConfig, db: LSMTree): Server = let registry = newDatabaseRegistry(config) let ctx = newExecutionContext(db, registry) @@ -389,7 +387,7 @@ proc handleClient(server: Server, client: AsyncSocket, clientId: int) {.async.} rest.add(more) let parts = rest.strip().split(" ") if parts.len >= 2: - let txnId = try: uint64(parseBiggestUint(parts[0])) except: 0'u64 + let txnId = try: uint64(parseBiggestUint(parts[0])) except CatchableError: 0'u64 let action = parts[1].toUpper() if server.distTxnManager != nil: let txn = server.distTxnManager.getTxn(txnId) @@ -428,8 +426,8 @@ proc handleClient(server: Server, client: AsyncSocket, clientId: int) {.async.} rest.add(more) let parts = rest.strip().split(" ") if parts.len >= 2: - let lsn = try: parseUInt(parts[0]) except: 0'u64 - let dataLen = try: parseInt(parts[1]) except: 0 + let lsn = try: parseUInt(parts[0]) except CatchableError: 0'u64 + let dataLen = try: parseInt(parts[1]) except CatchableError: 0 if dataLen > 0: var data = "" while data.len < dataLen: @@ -460,7 +458,7 @@ proc handleClient(server: Server, client: AsyncSocket, clientId: int) {.async.} let headerLine = "MIGRATE " & rest.strip() let parts = rest.strip().split(" ") if parts.len >= 2: - let entryCount = try: parseInt(parts[1]) except: 0 + let entryCount = try: parseInt(parts[1]) except CatchableError: 0 var data = "" if entryCount > 0: # Read all entries (each entry is key\0value\n) @@ -551,7 +549,7 @@ proc handleClient(server: Server, client: AsyncSocket, clientId: int) {.async.} # Shard-aware routing: check if this node should handle the write var shardCheck = true if server.clusterMembership.nodes.len > 0: - let stmts = try: parse(tokenize(queryStr)) except: nil + let stmts = try: parse(tokenize(queryStr)) except CatchableError: nil if stmts != nil: for stmt in stmts.stmts: if stmt.kind in {nkInsert, nkUpdate, nkDelete}: diff --git a/src/barabadb/core/sharding.nim b/src/barabadb/core/sharding.nim index fa06f10..17cbd73 100644 --- a/src/barabadb/core/sharding.nim +++ b/src/barabadb/core/sharding.nim @@ -337,8 +337,8 @@ proc handleMigrationMessage*(headerLine: string, data: string, if parts.len < 3: return "ERR invalid migrate header\n" - let shardId = try: parseInt(parts[1]) except: -1 - let entryCount = try: parseInt(parts[2]) except: 0 + let shardId = try: parseInt(parts[1]) except CatchableError: -1 + let entryCount = try: parseInt(parts[2]) except CatchableError: 0 if shardId < 0 or entryCount < 0: return "ERR invalid shard id or entry count\n" diff --git a/src/barabadb/core/tracing.nim b/src/barabadb/core/tracing.nim index 324c1e0..6838e96 100644 --- a/src/barabadb/core/tracing.nim +++ b/src/barabadb/core/tracing.nim @@ -128,5 +128,5 @@ proc exportOtlp*(tracer: Tracer, endpoint: string = "http://localhost:4318/v1/tr client.close() tracer.spans = @[] return true - except: + except CatchableError: return false diff --git a/src/barabadb/core/websocket.nim b/src/barabadb/core/websocket.nim index b025026..d95947a 100644 --- a/src/barabadb/core/websocket.nim +++ b/src/barabadb/core/websocket.nim @@ -189,7 +189,7 @@ proc notifyClient(client: WsClient, msg: string) {.async.} = try: let frame = encodeFrame(0x1, msg) await client.socket.send(frame) - except: + except CatchableError: discard proc broadcastToTable*(server: WsServer, table: string, msg: string) {.async.} = @@ -247,7 +247,7 @@ proc handleWsClient(server: WsServer, client: AsyncSocket, id: int) {.async.} = buf = buf[consumed..^1] - except: + except CatchableError: discard finally: echo "WebSocket client ", id, " disconnected" @@ -305,7 +305,7 @@ proc handleConnection(server: WsServer, client: AsyncSocket) {.async.} = await client.send("HTTP/1.1 401 Unauthorized\r\n\r\n") client.close() return - except: + except CatchableError: await client.send("HTTP/1.1 401 Unauthorized\r\n\r\n") client.close() return diff --git a/src/barabadb/protocol/zerocopy.nim b/src/barabadb/protocol/zerocopy.nim index 54071ab..e261c82 100644 --- a/src/barabadb/protocol/zerocopy.nim +++ b/src/barabadb/protocol/zerocopy.nim @@ -167,14 +167,14 @@ proc encodeRecord*(buf: var ZeroBuf, schema: ZcSchema, try: var v = int32(parseInt(value)) bigEndian32(addr buf.data[field.offset], unsafeAddr v) - except: + except CatchableError: var v: int32 = 0 bigEndian32(addr buf.data[field.offset], unsafeAddr v) of ztInt64: try: var v = int64(parseInt(value)) bigEndian64(addr buf.data[field.offset], unsafeAddr v) - except: + except CatchableError: var v: int64 = 0 bigEndian64(addr buf.data[field.offset], unsafeAddr v) of ztString: diff --git a/src/barabadb/query/executor.nim b/src/barabadb/query/executor.nim index fe98d11..e7726be 100644 --- a/src/barabadb/query/executor.nim +++ b/src/barabadb/query/executor.nim @@ -232,7 +232,7 @@ proc acquireMigrationLock(ctx: ExecutionContext): bool = let (locked, lockVal) = ctx.db.get(lockKey) if locked: # Check for stale lock (older than 1 hour) - let lockTime = try: parseInt(cast[string](lockVal)) except: 0 + let lockTime = try: parseInt(cast[string](lockVal)) except CatchableError: 0 if lockTime > 0 and (epochTime().int64 - lockTime) > 3600: # Stale lock — force release ctx.db.delete(lockKey) @@ -362,7 +362,7 @@ proc parseVectorString*(value: string): seq[float32] = if p.len > 0: try: result.add(parseFloat(p).float32) - except: + except CatchableError: discard # ---------------------------------------------------------------------- @@ -621,10 +621,10 @@ proc evalExpr*(expr: IRExpr, row: Row, ctx: ExecutionContext = nil): Value = case expr.valueKind of vkInt64: try: return Value(kind: vkInt64, int64Val: parseInt(s)) - except: return Value(kind: vkNull) + except CatchableError: return Value(kind: vkNull) of vkFloat64: try: return Value(kind: vkFloat64, float64Val: parseFloat(s)) - except: return Value(kind: vkNull) + except CatchableError: return Value(kind: vkNull) of vkBool: return Value(kind: vkBool, boolVal: s == "true") of vkNull: @@ -634,10 +634,10 @@ proc evalExpr*(expr: IRExpr, row: Row, ctx: ExecutionContext = nil): Value = if s.len == 0: return Value(kind: vkString, strVal: s) try: return Value(kind: vkInt64, int64Val: parseInt(s)) - except: + except CatchableError: try: return Value(kind: vkFloat64, float64Val: parseFloat(s)) - except: + except CatchableError: return Value(kind: vkString, strVal: s) return Value(kind: vkNull) of irekBinary: @@ -803,7 +803,7 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex of JNull: return "null" else: return $val return "" - except: + except CatchableError: return "" of irekBinary: let left = evalExprOld(expr.binLeft, row, ctx) @@ -814,30 +814,30 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex # Try numeric comparison try: if parseFloat(left) == parseFloat(right): return "true" - except: discard + except CatchableError: discard return "false" of irNeq: if left != right: return "true" # Try numeric comparison try: return if parseFloat(left) != parseFloat(right): "true" else: "false" - except: return "false" + except CatchableError: return "false" of irLt: try: return if parseFloat(left) < parseFloat(right): "true" else: "false" - except: return if left < right: "true" else: "false" + except CatchableError: return if left < right: "true" else: "false" of irLte: try: return if parseFloat(left) <= parseFloat(right): "true" else: "false" - except: return if left <= right: "true" else: "false" + except CatchableError: return if left <= right: "true" else: "false" of irGt: try: return if parseFloat(left) > parseFloat(right): "true" else: "false" - except: return if left > right: "true" else: "false" + except CatchableError: return if left > right: "true" else: "false" of irGte: try: return if parseFloat(left) >= parseFloat(right): "true" else: "false" - except: return if left >= right: "true" else: "false" + except CatchableError: return if left >= right: "true" else: "false" of irAnd: if left == "true" and right == "true": return "true" return "false" @@ -867,7 +867,7 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex try: let rePattern = re(pattern) if left.match(rePattern): return "true" - except: discard + except CatchableError: discard return "false" of irILike: proc escapeRe(s: string): string = @@ -882,7 +882,7 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex try: let rePattern = re(pattern) if left.toLower().match(rePattern): return "true" - except: discard + except CatchableError: discard return "false" of irIn: if expr.binRight.kind == irekSubquery: @@ -902,7 +902,7 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex let lv = parseFloat(left) let rv = parseFloat(right) return if lv == rv: "true" else: "false" - except: discard + except CatchableError: discard return if left == right: "true" else: "false" of irNotIn: if expr.binRight.kind == irekSubquery: @@ -922,7 +922,7 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex let lv = parseFloat(left) let rv = parseFloat(right) return if lv != rv: "true" else: "false" - except: discard + except CatchableError: discard return if left != right: "true" else: "false" of irFtsMatch: # Check for FTS index via ctx @@ -986,7 +986,7 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex return "true" else: return if $(leftNode) == $(rightNode): "true" else: "false" - except: + except CatchableError: return "false" of irJsonContainedBy: # Check if left JSON is contained by right JSON (reverse of contains) @@ -1010,7 +1010,7 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex return "true" else: return if $(leftNode) == $(rightNode): "true" else: "false" - except: + except CatchableError: return "false" of irJsonHasAny: # Check if JSON object has any of the keys in right array @@ -1022,7 +1022,7 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex if key.kind == JString and leftNode.hasKey(key.getStr()): return "true" return "false" - except: + except CatchableError: return "false" of irJsonHasAll: # Check if JSON object has all of the keys in right array @@ -1035,7 +1035,7 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex return "false" return "true" return "false" - except: + except CatchableError: return "false" else: return "false" of irekUnary: @@ -1097,10 +1097,10 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex try: let idx = parseInt(key) return if idx >= 0 and idx < node.len: "true" else: "false" - except: + except CatchableError: return "false" return "false" - except: + except CatchableError: return "false" of "current_setting": if expr.irFuncArgs.len < 1: @@ -1460,13 +1460,13 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex try: let dt = parse(val, "yyyy-MM-dd HH:mm:ss") return $(dt.toTime().toUnix()) - except: + except CatchableError: return "0" elif fmt == "%Y-%m-%dT%H:%M:%SZ": try: let dt = parse(val, "yyyy-MM-dd HH:mm:ss") return format(dt, "yyyy-MM-dd'T'HH:mm:ss'Z'") - except: + except CatchableError: return "" return "" else: @@ -1773,7 +1773,7 @@ proc execInsert*(ctx: ExecutionContext, table: string, fields: seq[string], valu props[f] = rowVals[i] try: gengine.addNodeWithId(graph, nid, label, props) - except: + except CatchableError: discard elif table == graphName & "_edges": var srcStr = "" @@ -1786,13 +1786,13 @@ proc execInsert*(ctx: ExecutionContext, table: string, fields: seq[string], valu elif f == "dest_id": dstStr = rowVals[i] elif f == "edge_label": label = rowVals[i] elif f == "weight": - try: weight = parseFloat(rowVals[i]) except: discard + try: weight = parseFloat(rowVals[i]) except CatchableError: discard if srcStr.len > 0 and dstStr.len > 0: let srcId = gengine.NodeId(parseUInt(srcStr)) let dstId = gengine.NodeId(parseUInt(dstStr)) try: gengine.addEdgeWithId(graph, srcId, dstId, label, weight) - except: + except CatchableError: discard inc count @@ -2018,10 +2018,10 @@ proc validateType*(colType: string, value: string): (bool, string) = let t = colType.toUpper() if t == "INTEGER" or t == "INT" or t == "BIGINT" or t == "SMALLINT" or t == "SERIAL": try: discard parseInt(value) - except: return (false, "Type mismatch: expected " & t & " but got '" & value & "'") + except CatchableError: return (false, "Type mismatch: expected " & t & " but got '" & value & "'") elif t == "FLOAT" or t == "REAL" or t == "DOUBLE" or t == "DOUBLE PRECISION" or t == "NUMERIC": try: discard parseFloat(value) - except: return (false, "Type mismatch: expected " & t & " but got '" & value & "'") + except CatchableError: return (false, "Type mismatch: expected " & t & " but got '" & value & "'") elif t == "BOOLEAN" or t == "BOOL": let lv = value.toLower() if lv notin ["true", "false", "1", "0", "t", "f", "yes", "no"]: @@ -2032,7 +2032,7 @@ proc validateType*(colType: string, value: string): (bool, string) = elif t == "JSON" or t == "JSONB": try: discard parseJson(value) - except: + except CatchableError: return (false, "Type mismatch: expected JSON but got '" & value & "'") elif t.startsWith("VECTOR"): let vec = parseVectorString(value) @@ -2044,7 +2044,7 @@ proc validateType*(colType: string, value: string): (bool, string) = if dimStart >= 0 and dimEnd > dimStart: try: expectedDim = parseInt(t[dimStart+1.. 0 and vec.len != expectedDim: return (false, "Vector dimension mismatch: expected " & $expectedDim & " but got " & $vec.len) @@ -2572,7 +2572,7 @@ proc compareRowsByOrder(a, b: Row, orderExprs: seq[IRExpr], orderDirs: seq[bool] let fb = parseFloat(valueToString(vb)) if fa < fb: cmpRes = -1 elif fa > fb: cmpRes = 1 - except: + except CatchableError: cmpRes = cmp(valueToString(va), valueToString(vb)) if cmpRes != 0: return if orderDirs.len > i and orderDirs[i]: -cmpRes else: cmpRes @@ -2592,12 +2592,12 @@ proc resolveFrameBounds(pos, partLen: int, frameStart, frameEnd: string): (int, elif frameStart.endsWith(" PRECEDING"): let nStr = frameStart[0..^11] var n = 0 - try: n = parseInt(nStr) except: n = 0 + try: n = parseInt(nStr) except CatchableError: n = 0 startPos = max(0, pos - n) elif frameStart.endsWith(" FOLLOWING"): let nStr = frameStart[0..^11] var n = 0 - try: n = parseInt(nStr) except: n = 0 + try: n = parseInt(nStr) except CatchableError: n = 0 startPos = min(partLen - 1, pos + n) # Parse end boundary @@ -2608,12 +2608,12 @@ proc resolveFrameBounds(pos, partLen: int, frameStart, frameEnd: string): (int, elif frameEnd.endsWith(" PRECEDING"): let nStr = frameEnd[0..^11] var n = 0 - try: n = parseInt(nStr) except: n = 0 + try: n = parseInt(nStr) except CatchableError: n = 0 endPos = max(0, pos - n) elif frameEnd.endsWith(" FOLLOWING"): let nStr = frameEnd[0..^11] var n = 0 - try: n = parseInt(nStr) except: n = 0 + try: n = parseInt(nStr) except CatchableError: n = 0 endPos = min(partLen - 1, pos + n) if startPos > endPos: @@ -2668,7 +2668,7 @@ proc computeWindowValues*(rows: seq[Row], expr: IRExpr, ctx: ExecutionContext = of "ntile": var n = 1 if expr.wfArgs.len > 0: - try: n = parseInt(valueToString(evalExpr(expr.wfArgs[0], rows[sortedIdxs[0]], ctx))) except: n = 1 + try: n = parseInt(valueToString(evalExpr(expr.wfArgs[0], rows[sortedIdxs[0]], ctx))) except CatchableError: n = 1 if n < 1: n = 1 let groupSize = sortedIdxs.len div n let remainder = sortedIdxs.len mod n @@ -2687,7 +2687,7 @@ proc computeWindowValues*(rows: seq[Row], expr: IRExpr, ctx: ExecutionContext = var offset = 1 var defaultVal = "" if expr.wfArgs.len > 1: - try: offset = parseInt(valueToString(evalExpr(expr.wfArgs[1], rows[sortedIdxs[0]], ctx))) except: offset = 1 + try: offset = parseInt(valueToString(evalExpr(expr.wfArgs[1], rows[sortedIdxs[0]], ctx))) except CatchableError: offset = 1 if expr.wfArgs.len > 2: defaultVal = valueToString(evalExpr(expr.wfArgs[2], rows[sortedIdxs[0]], ctx)) for pos, rowIdx in sortedIdxs: @@ -2700,7 +2700,7 @@ proc computeWindowValues*(rows: seq[Row], expr: IRExpr, ctx: ExecutionContext = var offset = 1 var defaultVal = "" if expr.wfArgs.len > 1: - try: offset = parseInt(valueToString(evalExpr(expr.wfArgs[1], rows[sortedIdxs[0]], ctx))) except: offset = 1 + try: offset = parseInt(valueToString(evalExpr(expr.wfArgs[1], rows[sortedIdxs[0]], ctx))) except CatchableError: offset = 1 if expr.wfArgs.len > 2: defaultVal = valueToString(evalExpr(expr.wfArgs[2], rows[sortedIdxs[0]], ctx)) for pos, rowIdx in sortedIdxs: @@ -2843,14 +2843,14 @@ proc executePlan*(ctx: ExecutionContext, plan: IRPlan): seq[Row] = var sum = 0.0 for row in filteredRows: let v = evalExpr(expr.aggArgs[0], row, ctx) - try: sum += parseFloat(valueToString(v)) except: discard + try: sum += parseFloat(valueToString(v)) except CatchableError: discard newRow[alias] = $sum of irAvg: var sum = 0.0 var count = 0 for row in filteredRows: let v = evalExpr(expr.aggArgs[0], row, ctx) - try: sum += parseFloat(valueToString(v)); count += 1 except: discard + try: sum += parseFloat(valueToString(v)); count += 1 except CatchableError: discard newRow[alias] = if count > 0: $(sum / float(count)) else: "0" of irMin: var minVal = "" @@ -2930,7 +2930,7 @@ proc executePlan*(ctx: ExecutionContext, plan: IRPlan): seq[Row] = let fb = parseFloat(valueToString(vb)) if fa < fb: cmpRes = -1 elif fa > fb: cmpRes = 1 - except: + except CatchableError: cmpRes = cmp(valueToString(va), valueToString(vb)) if not ascending: cmpRes = -cmpRes if cmpRes != 0: return cmpRes @@ -3022,14 +3022,14 @@ proc executePlan*(ctx: ExecutionContext, plan: IRPlan): seq[Row] = var sum = 0.0 for row in filteredRows: let v = evalExpr(aggExpr.aggArgs[0], row, ctx) - try: sum += parseFloat(valueToString(v)) except: discard + try: sum += parseFloat(valueToString(v)) except CatchableError: discard aggRow[aggKey] = $sum of irAvg: var sum = 0.0 var count = 0 for row in filteredRows: let v = evalExpr(aggExpr.aggArgs[0], row, ctx) - try: sum += parseFloat(valueToString(v)); count += 1 except: discard + try: sum += parseFloat(valueToString(v)); count += 1 except CatchableError: discard aggRow[aggKey] = if count > 0: $(sum / float(count)) else: "0" of irMin: var minVal = "" @@ -3544,14 +3544,14 @@ proc executePlan*(ctx: ExecutionContext, plan: IRPlan): seq[Row] = var sum = 0.0 for row in matchingRows: let v = evalExpr(plan.pivotAgg.aggArgs[0], row, ctx) - try: sum += parseFloat(valueToString(v)) except: discard + try: sum += parseFloat(valueToString(v)) except CatchableError: discard aggResult = $sum of irAvg: var sum = 0.0 var count = 0 for row in matchingRows: let v = evalExpr(plan.pivotAgg.aggArgs[0], row, ctx) - try: sum += parseFloat(valueToString(v)); count += 1 except: discard + try: sum += parseFloat(valueToString(v)); count += 1 except CatchableError: discard aggResult = if count > 0: $(sum / float(count)) else: "0" of irMin: var minVal = "" @@ -3609,8 +3609,8 @@ proc executePlan*(ctx: ExecutionContext, plan: IRPlan): seq[Row] = let algo = plan.graphAlgo.toLowerAscii() let returnCols = plan.graphReturnCols let firstNodeId = if g.nodes.len > 0: g.nodes.keys.toSeq[0] else: gengine.NodeId(0) - let explicitStart = try: parseUInt(plan.graphStartNode) except: 0'u64 - let explicitEnd = try: parseUInt(plan.graphEndNode) except: 0'u64 + let explicitStart = try: parseUInt(plan.graphStartNode) except CatchableError: 0'u64 + let explicitEnd = try: parseUInt(plan.graphEndNode) except CatchableError: 0'u64 case algo of "bfs": @@ -4034,7 +4034,7 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu if fa < fb: return -1 if fa > fb: return 1 return 0 - except: + except CatchableError: return cmp(valueToString(va), valueToString(vb)) filteredRows.sort(sortCmp, if asc: Ascending else: Descending) if stmt.selLimit != nil: @@ -4369,7 +4369,7 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu ctx.autoIncCounters[counterKey] = intVal + 1 finally: release(ctx.sharedLock.lock) - except: discard + except CatchableError: discard applyDefaultValues(tbl, mutableFields, mutableValues) diff --git a/src/barabadb/query/udf.nim b/src/barabadb/query/udf.nim index a169ad6..a4ccce7 100644 --- a/src/barabadb/query/udf.nim +++ b/src/barabadb/query/udf.nim @@ -211,7 +211,7 @@ proc registerStdlib*(reg: UDFRegistry) = if args.len > 0 and args[0].kind == vkString: try: return Value(kind: vkInt64, int64Val: parseInt(args[0].strVal)) - except: + except CatchableError: discard return Value(kind: vkNull)) diff --git a/src/barabadb/storage/compaction.nim b/src/barabadb/storage/compaction.nim index 4dd2bdd..9404c1e 100644 --- a/src/barabadb/storage/compaction.nim +++ b/src/barabadb/storage/compaction.nim @@ -58,7 +58,7 @@ proc rebuildFromLSM*(cs: CompactionStrategy, db: LSMTree) = cs.clear() cs.dataDir = db.dir for sst in db.sstables: - let size = try: int(getFileSize(sst.path)) except: sst.entryCount * 64 + let size = try: int(getFileSize(sst.path)) except CatchableError: sst.entryCount * 64 cs.addTable(SSTableMeta( path: sst.path, level: sst.level, @@ -143,7 +143,7 @@ proc compact*(cs: CompactionStrategy, level: int): CompactionResult = var sst = writeSSTable(final, outputPath, level + 1) # Use actual file size instead of rough guess - let actualSize = try: getFileSize(outputPath) except: final.len * 64 + let actualSize = try: getFileSize(outputPath) except CatchableError: final.len * 64 let outputMeta = SSTableMeta( path: outputPath, @@ -159,7 +159,7 @@ proc compact*(cs: CompactionStrategy, level: int): CompactionResult = let (ok, msg) = verifySSTable(outputPath) if not ok: echo "[ERROR] Compaction output verification failed: ", msg - try: removeFile(outputPath) except: discard + try: removeFile(outputPath) except CatchableError: discard return CompactionResult() # Remove old SSTable files diff --git a/src/barabadb/storage/gate.nim b/src/barabadb/storage/gate.nim index a2d7bc6..badbd25 100644 --- a/src/barabadb/storage/gate.nim +++ b/src/barabadb/storage/gate.nim @@ -18,13 +18,15 @@ var proc initStorageGate*() = ## Idempotent when called from a single thread at startup. + ## Must run before multi-threaded accept (HTTP workers / TCP). if not gInited: initLock(gGate) gInited = true proc acquireStorageGate*() {.inline.} = ## Prefer calling initStorageGate() once at process start (main). - ## Lazy-init is allowed for unit tests (single-threaded). + ## Lazy-init is only safe for single-threaded unit tests — concurrent + ## first-time init races on gInited / initLock. if not gInited: initStorageGate() acquire(gGate) diff --git a/src/barabadb/storage/lsm.nim b/src/barabadb/storage/lsm.nim index 4327aee..c310a88 100644 --- a/src/barabadb/storage/lsm.nim +++ b/src/barabadb/storage/lsm.nim @@ -473,7 +473,7 @@ proc listLegacySSTables*(dir: string): seq[(string, uint32)] = let sst = loadSSTable(path) if sst.fileVersion < SSTableVersion: result.add((path, sst.fileVersion)) - except: + except CatchableError: discard proc migrateSSTable*(path: string): bool = @@ -595,7 +595,7 @@ proc checkStorageConsistency*(db: LSMTree): seq[string] = let j = parseJson(readFile(manifestPath)) for node in j{"sstables"}: manifestPaths.add(node{"path"}.getStr()) - except: + except CatchableError: result.add("MANIFEST is corrupt or unreadable") return diff --git a/src/barabadb/storage/wal.nim b/src/barabadb/storage/wal.nim index d4e4fc9..1d4cab2 100644 --- a/src/barabadb/storage/wal.nim +++ b/src/barabadb/storage/wal.nim @@ -77,7 +77,7 @@ proc parseWalSequence*(filename: string): int64 = result = parseBiggestInt(numStr) else: result = 0 - except: + except CatchableError: result = 0 proc listWalArchive*(dir: string): seq[WalSegment] = @@ -90,7 +90,7 @@ proc listWalArchive*(dir: string): seq[WalSegment] = if kind == pcFile and path.endsWith(".log"): let seqNum = parseWalSequence(extractFilename(path)) if seqNum > 0: - let size = try: getFileSize(path) except: 0 + let size = try: getFileSize(path) except CatchableError: 0 result.add(WalSegment(sequence: seqNum, path: path, size: size)) result.sort(proc(a, b: WalSegment): int = cmp(a.sequence, b.sequence)) @@ -140,7 +140,7 @@ proc maybeRotate*(wal: var WriteAheadLog) = ## Rotate if current WAL exceeds max segment size. if wal.maxSegmentSize <= 0: return - let currentSize = try: getFileSize(wal.path) except: 0 + let currentSize = try: getFileSize(wal.path) except CatchableError: 0 if currentSize >= wal.maxSegmentSize: wal.rotate() diff --git a/src/baradadb.nim b/src/baradadb.nim index aa2a420..5eb32ba 100644 --- a/src/baradadb.nim +++ b/src/baradadb.nim @@ -61,7 +61,7 @@ proc applyCompactionResult(db: LSMTree, result: compaction.CompactionResult) = var sst = loadSSTable(meta.path) let name = splitFile(meta.path).name # Prefer numeric id from filename; otherwise allocate - let parsed = try: parseInt(name) except: -1 + let parsed = try: parseInt(name) except CatchableError: -1 if parsed >= 0: sst.id = parsed else: @@ -363,7 +363,7 @@ proc main() = let parts = seed.strip().split(":") if parts.len >= 2: let host = parts[0] - let port = try: parseInt(parts[1]) except: 0 + let port = try: parseInt(parts[1]) except CatchableError: 0 if port > 0: let seedNode = newGossipNode(host & ":" & $port, host, port) tcpServer.gossipProtocol.join(seedNode) diff --git a/src/baramcp.nim b/src/baramcp.nim index bee4809..8d8ac6b 100644 --- a/src/baramcp.nim +++ b/src/baramcp.nim @@ -18,7 +18,7 @@ when isMainModule: try: discard server.init(dataDir) server.run() - except: + except CatchableError: server.logToStderr("Fatal error: " & getCurrentExceptionMsg()) finally: server.close() diff --git a/tests/config.nims b/tests/config.nims index a119208..16c2eb5 100644 --- a/tests/config.nims +++ b/tests/config.nims @@ -1 +1,2 @@ --path:"../src" +--path:"../clients/nim/src" diff --git a/tests/test_all.nim b/tests/test_all.nim index b4514dc..273c8a0 100644 --- a/tests/test_all.nim +++ b/tests/test_all.nim @@ -31,7 +31,11 @@ import barabadb/query/udf import barabadb/vector/simd import barabadb/core/crossmodal import barabadb/core/gossip -import barabadb/client/client +# Canonical Nim client (not the deprecated src/barabadb/client). +# Selective import avoids WireValue clash with barabadb/protocol/wire. +from baradb/client import parseConnectionString, defaultClientConfig, newBaraClient, + newQueryBuilder, newSyncClient, BaraClient, QueryBuilder, ClientConfig, SyncClient, + select, `from`, where, join, leftJoin, groupBy, having, orderBy, limit, offset, build, exec import barabadb/client/fileops import barabadb/fts/multilang as mlang import barabadb/protocol/zerocopy diff --git a/tests/test_schema_persist.nim b/tests/test_schema_persist.nim index e655bae..178dff0 100644 --- a/tests/test_schema_persist.nim +++ b/tests/test_schema_persist.nim @@ -6,7 +6,6 @@ import std/tables import barabadb/storage/lsm import barabadb/query/executor import barabadb/query/parser -import barabadb/query/ast proc execSql(ctx: ExecutionContext, sql: string): ExecResult = let node = parse(sql) diff --git a/tests/test_storage_hardening.nim b/tests/test_storage_hardening.nim index 48155a6..6ed8903 100644 --- a/tests/test_storage_hardening.nim +++ b/tests/test_storage_hardening.nim @@ -1,7 +1,6 @@ ## Focused storage hardening tests (avoids full suite compile issues) import std/unittest import std/os -import std/strutils import std/locks import barabadb/storage/lsm import barabadb/storage/rwlock