Source code for falkordb.falkordb

from typing import List, Optional, Union

import redis  # type: ignore[import-not-found]
from redis.driver_info import DriverInfo
from redis.exceptions import RedisError

from ._version import get_package_version
from .cluster import Cluster_Conn, Is_Cluster
from .graph import Graph
from .sentinel import Is_Sentinel, Sentinel_Conn

# config commands
UDF_CMD = "GRAPH.UDF"
LIST_CMD = "GRAPH.LIST"
CONFIG_CMD = "GRAPH.CONFIG"


[docs] class FalkorDB: """ FalkorDB Class for interacting with a FalkorDB server. Supports both TCP and Unix socket connections. Usage example:: from falkordb import FalkorDB # connect to the database and select the 'social' graph db = FalkorDB() graph = db.select_graph("social") # connect using a Unix socket db = FalkorDB.from_url("unix:///path/to/redis.sock") graph = db.select_graph("social") # get a single 'Person' node from the graph and print its name result = graph.query("MATCH (n:Person) RETURN n LIMIT 1").result_set person = result[0][0] print(person.properties['name']) """ def __init__( self, host="localhost", port=6379, password=None, socket_timeout=None, socket_connect_timeout=None, socket_keepalive=None, socket_keepalive_options=None, connection_pool=None, unix_socket_path=None, encoding="utf-8", encoding_errors="strict", retry_on_error=None, ssl=False, ssl_keyfile=None, ssl_certfile=None, ssl_cert_reqs="required", ssl_ca_certs=None, ssl_ca_path=None, ssl_ca_data=None, ssl_check_hostname=False, ssl_password=None, ssl_validate_ocsp=False, ssl_validate_ocsp_stapled=False, ssl_ocsp_context=None, ssl_ocsp_expected_cert=None, max_connections=None, single_connection_client=False, health_check_interval=0, client_name=None, lib_name="FalkorDB", lib_version=None, username=None, retry=None, connect_func=None, credential_provider=None, protocol=2, # FalkorDB Cluster Params cluster_error_retry_attempts=3, startup_nodes=None, require_full_coverage=False, reinitialize_steps=5, read_from_replicas=False, dynamic_startup_nodes=True, url=None, address_remap=None, ): conn = redis.Redis( host=host, port=port, db=0, password=password, socket_timeout=socket_timeout, socket_connect_timeout=socket_connect_timeout, socket_keepalive=socket_keepalive, socket_keepalive_options=socket_keepalive_options, connection_pool=connection_pool, unix_socket_path=unix_socket_path, encoding=encoding, encoding_errors=encoding_errors, decode_responses=True, retry_on_error=retry_on_error, ssl=ssl, ssl_keyfile=ssl_keyfile, ssl_certfile=ssl_certfile, ssl_cert_reqs=ssl_cert_reqs, ssl_ca_certs=ssl_ca_certs, ssl_ca_path=ssl_ca_path, ssl_ca_data=ssl_ca_data, ssl_check_hostname=ssl_check_hostname, ssl_password=ssl_password, ssl_validate_ocsp=ssl_validate_ocsp, ssl_validate_ocsp_stapled=ssl_validate_ocsp_stapled, ssl_ocsp_context=ssl_ocsp_context, ssl_ocsp_expected_cert=ssl_ocsp_expected_cert, max_connections=max_connections, single_connection_client=single_connection_client, health_check_interval=health_check_interval, client_name=client_name, driver_info=DriverInfo( lib_name, get_package_version() if lib_version is None else lib_version, ), username=username, retry=retry, redis_connect_func=connect_func, credential_provider=credential_provider, protocol=protocol, ) if Is_Sentinel(conn): self.sentinel, self.service_name = Sentinel_Conn(conn, ssl) conn = self.sentinel.master_for(self.service_name, ssl=ssl) if Is_Cluster(conn): conn = Cluster_Conn( conn, ssl, cluster_error_retry_attempts, startup_nodes, require_full_coverage, reinitialize_steps, read_from_replicas, dynamic_startup_nodes, url, address_remap, ) self.connection = conn self.flushdb = conn.flushdb self.execute_command = conn.execute_command
[docs] @classmethod def from_url(cls, url: str, **kwargs) -> "FalkorDB": """ Creates a new FalkorDB instance from a URL. Args: cls: The class itself. url (str): The URL. kwargs: Additional keyword arguments to pass to the ``DB.from_url`` function. Returns: DB: A new DB instance. Usage example:: db = FalkorDB.from_url("falkor://[[username]:[password]]@localhost:6379") db = FalkorDB.from_url("falkors://[[username]:[password]]@localhost:6379") db = FalkorDB.from_url("unix://[username@]/path/to/socket.sock?db=0[&password=password]") """ # switch from redis:// to falkordb:// if url.startswith("falkor://"): url = "redis://" + url[len("falkor://") :] elif url.startswith("falkors://"): url = "rediss://" + url[len("falkors://") :] kwargs["decode_responses"] = True conn = redis.from_url(url, **kwargs) return cls(connection_pool=conn.connection_pool)
[docs] def select_graph(self, graph_id: str) -> Graph: """ Selects a graph by creating a new Graph instance. Args: graph_id (str): The identifier of the graph. Returns: Graph: A new Graph instance associated with the selected graph. """ if not isinstance(graph_id, str) or graph_id == "": raise TypeError( f"Expected a string parameter, but received {type(graph_id)}." ) return Graph(self, graph_id)
[docs] def list_graphs(self) -> List[str]: """ Lists all graph names. See: https://docs.falkordb.com/commands/graph.list.html Returns: List: List of graph names. """ return self.connection.execute_command(LIST_CMD)
[docs] def config_get(self, name: str) -> Union[int, str]: """ Retrieve a DB level configuration. For a list of available configurations see: https://docs.falkordb.com/configuration.html#falkordb-configuration-parameters Args: name (str): The name of the configuration. Returns: int or str: The configuration value. """ return self.connection.execute_command(CONFIG_CMD, "GET", name)[1]
[docs] def config_set(self, name: str, value=None) -> None: """ Update a DB level configuration. For a list of available configurations see: https://docs.falkordb.com/configuration.html#falkordb-configuration-parameters Args: name (str): The name of the configuration. value: The value to set. Returns: None """ return self.connection.execute_command(CONFIG_CMD, "SET", name, value)
[docs] def close(self) -> None: """ Close the underlying connection(s). """ try: self.connection.close() except RedisError: # best-effort close — don't raise on Redis errors pass
def __enter__(self) -> "FalkorDB": """Return self to support usage in a with-statement.""" return self def __exit__(self, exc_type, exc_val, exc_tb) -> None: """Close the connection when exiting a with-statement.""" self.close() # GRAPH.UDF LOAD [REPLACE] <lib> <script>
[docs] def udf_load(self, name: str, script: str, replace: bool = False): """ Load a User Defined Function (UDF) library. Args: name (str): The name of the library to load. script (str): The UDF script contents. replace (bool, optional): If True, replace an existing library with the same name. Defaults to False. """ # prep arguments args = [UDF_CMD, "LOAD"] if replace: args.append("REPLACE") args.extend([name, script]) # propagate command in cluster mode if Is_Cluster(self.connection): for node in self.connection.get_primaries(): # create a direct connection to this node client = self.connection.get_redis_connection(node) resp = client.execute_command(*args) else: resp = self.connection.execute_command(*args) return resp
# GRAPH.UDF LIST [LIBRARYNAME] [WITHCODE]
[docs] def udf_list(self, lib: Optional[str] = None, with_code: bool = False): """ List User Defined Function (UDF) libraries. Args: lib (str, optional): If provided, filter the list to this specific library. with_code (bool, optional): If True, include the library source code in the result. Defaults to False. Returns: list: A list of UDF libraries and their metadata. """ args = [UDF_CMD, "LIST"] if lib is not None: args.append(lib) if with_code: args.append("WITHCODE") return self.connection.execute_command(*args)
# GRAPH.UDF FLUSH
[docs] def udf_flush(self): """ Flush (remove) all User Defined Function (UDF) libraries. """ # propagate command in cluster mode if Is_Cluster(self.connection): for node in self.connection.get_primaries(): # create a direct connection to this node client = self.connection.get_redis_connection(node) resp = client.execute_command(UDF_CMD, "FLUSH") else: resp = self.connection.execute_command(UDF_CMD, "FLUSH") return resp
# GRAPH.UDF DELETE <lib>
[docs] def udf_delete(self, lib: str): """ Delete a User Defined Function (UDF) library. Args: lib (str): The name of the library to delete. """ # propagate command in cluster mode if Is_Cluster(self.connection): for node in self.connection.get_primaries(): # create a direct connection to this node client = self.connection.get_redis_connection(node) resp = client.execute_command(UDF_CMD, "DELETE", lib) else: resp = self.connection.execute_command(UDF_CMD, "DELETE", lib) return resp