# pylint: disable=W0718,R0904
"""
Smarter API Vectorstore Manifest handler.
The broker converts between Vectorstore manifests and the
:class:`~smarter.apps.vectorstore.models.VectorstoreMeta` model, and implements the ``smarter``
CLI commands for Vectorstores:
- ``apply``: create or update the Vectorstore. It does not create the database.
- ``deploy``: create the database: a self-hosted Qdrant server, or a managed index.
- ``undeploy``: stop serving. A self-hosted server's data is kept.
- ``delete``: destroy the database and all of its data, unless deletionProtection is enabled,
then delete the Vectorstore.
- ``describe``, ``get``, ``logs`` and ``example_manifest``.
Account admins, i.e. staff, may apply, deploy, undeploy and delete their own Vectorstores.
Anyone may describe and get the Vectorstores that are shared with them.
"""
import datetime
from typing import Any, Optional, Type
from django.db import transaction
from django.http import HttpRequest
from rest_framework.serializers import ModelSerializer
from smarter.apps.account.utils import smarter_cached_objects
from smarter.apps.connection.models import ApiConnection
from smarter.apps.plugin.signals import broker_ready
from smarter.apps.provider.models import Provider
from smarter.apps.vectorstore.caching import (
invalidate_all_cached_vectorstores_for_user_profile,
)
from smarter.apps.vectorstore.manifest.models.vectorstore.const import MANIFEST_KIND
from smarter.apps.vectorstore.manifest.models.vectorstore.metadata import (
SAMVectorstoreMetadata,
)
from smarter.apps.vectorstore.manifest.models.vectorstore.model import SAMVectorstore
from smarter.apps.vectorstore.manifest.models.vectorstore.spec import (
SAMVectorstoreEmbeddings,
SAMVectorstoreIndex,
SAMVectorstoreMaintenance,
SAMVectorstoreSelfHosted,
SAMVectorstoreSpec,
)
from smarter.apps.vectorstore.manifest.models.vectorstore.status import (
SAMVectorstoreStatus,
)
from smarter.apps.vectorstore.models import (
VectorstoreDocumentStatus,
VectorstoreMeta,
VectorstoreStatus,
)
from smarter.apps.vectorstore.serializers import VectorstoreSerializer
from smarter.apps.vectorstore.service import VectorstoreService
from smarter.lib import logging
from smarter.lib.django.waffle import SmarterWaffleSwitches
from smarter.lib.journal.enum import SmarterJournalCliCommands
from smarter.lib.journal.http import SmarterJournaledJsonResponse
from smarter.lib.manifest.broker import (
AbstractBroker,
SAMBrokerError,
SAMBrokerErrorNotFound,
SAMBrokerErrorNotImplemented,
SAMBrokerErrorNotReady,
)
from smarter.lib.manifest.enum import (
SAMKeys,
SAMMetadataKeys,
SCLIResponseGet,
SCLIResponseGetData,
)
logger = logging.getSmarterLogger(
__name__, any_switches=[SmarterWaffleSwitches.VECTORSTORE_LOGGING, SmarterWaffleSwitches.MANIFEST_LOGGING]
)
MAX_RESULTS = 1000
IMMUTABLE_WHILE_DEPLOYED = {
# what cannot change once the database exists, because it would no longer match it.
"backend": lambda spec: spec.backend,
"hosting": lambda spec: spec.hosting,
"index.name": lambda spec: spec.index.name,
"index.dimension": lambda spec: spec.index.dimension,
"index.metric": lambda spec: spec.index.metric,
"selfHosted.storage": lambda spec: spec.selfHosted.storage if spec.selfHosted else None,
"pinecone": lambda spec: spec.pinecone.model_dump() if spec.pinecone else None,
}
[docs]
class SAMVectorstoreBrokerError(SAMBrokerError):
"""Base exception for Smarter API Vectorstore Broker handling."""
@property
def get_formatted_err_message(self):
return "Smarter API Vectorstore Manifest Broker Error"
[docs]
class SAMVectorstoreBroker(AbstractBroker):
"""Broker for Vectorstore manifests.
See the module's documentation.
"""
_manifest: Optional[SAMVectorstore] = None
_pydantic_model: Type[SAMVectorstore] = SAMVectorstore
_vectorstore: Optional[VectorstoreMeta] = None
_name: Optional[str] = None
_ready: bool = False
[docs]
def __init__(self, *args, **kwargs) -> None:
super().__init__(*args, **kwargs)
logger.info(
"%s.__init__() broker for %s %s is %s.", self.formatted_class_name, self.kind, self.name, self.ready_state
)
@property
def SerializerClass(self) -> Type[ModelSerializer]:
return VectorstoreSerializer
@property
def ready(self) -> bool:
"""A broker is ready if it has a manifest, or an account."""
if self._ready:
return self._ready
if not super().ready:
return False
if self.manifest is not None or self.account is not None:
self._ready = True
broker_ready.send(sender=self.__class__, broker=self)
return self._ready
@property
def vectorstore(self) -> Optional[VectorstoreMeta]:
"""The user's own Vectorstore with the broker's name, else one shared with them.
It is never created here.
"""
if self._vectorstore:
return self._vectorstore
if not self.user_profile or not self.name:
return None
self._vectorstore = (
VectorstoreMeta.objects.filter(user_profile=self.user_profile, name=self.name).first()
or VectorstoreMeta.objects.filter(name=self.name)
.with_read_permission_for(self.user_profile.user) # type: ignore[attr-defined]
.order_by("-updated_at")
.first()
)
return self._vectorstore
[docs]
def owned_vectorstore(self, command: SmarterJournalCliCommands) -> VectorstoreMeta:
"""The Vectorstore, if the user is staff and owns it."""
if not self.user_profile:
raise SAMBrokerErrorNotReady("user_profile is not set.", thing=self.kind, command=command)
if not (self.user_profile.user.is_staff or self.user_profile.user.is_superuser):
raise SAMVectorstoreBrokerError(
f"Only account admins may {command.value} a {self.kind}.", thing=self.kind, command=command
)
vectorstore = self.vectorstore
if vectorstore is None or (
vectorstore.user_profile_id != self.user_profile.pk # type: ignore[attr-defined]
and not self.user_profile.user.is_superuser
):
raise SAMBrokerErrorNotFound(f"{self.kind} {self.name} not found", thing=self.kind, command=command)
return vectorstore
# -------------------------------------------------------------------------
# resolving names
# -------------------------------------------------------------------------
[docs]
def resolve_provider(self, name: str) -> Provider:
"""Spec.embeddings.provider: the user's own Provider, else the most recently updated one shared with them."""
assert self.user_profile is not None
provider = Provider.objects.filter(name=name, user_profile=self.user_profile).first() or (
Provider.objects.filter(name=name)
.with_read_permission_for(self.user_profile.user) # type: ignore[attr-defined]
.order_by("-updated_at")
.first()
)
if provider is None:
raise SAMBrokerErrorNotFound(
f"spec.embeddings.provider: Provider {name} not found, or not shared with you.",
thing=self.kind,
command=SmarterJournalCliCommands.APPLY,
)
return provider
[docs]
def resolve_connection(self, name: Optional[str]) -> Optional[ApiConnection]:
"""Spec.connection: the user's own ApiConnection, else one shared with them."""
if not name:
return None
assert self.user_profile is not None
connection = ApiConnection.objects.filter(name=name, user_profile=self.user_profile).first() or (
ApiConnection.objects.filter(name=name)
.with_read_permission_for(self.user_profile.user) # type: ignore[attr-defined]
.order_by("-updated_at")
.first()
)
if connection is None:
raise SAMBrokerErrorNotFound(
f"spec.connection: ApiConnection {name} not found, or not shared with you.",
thing=self.kind,
command=SmarterJournalCliCommands.APPLY,
)
return connection
# -------------------------------------------------------------------------
# conversions
# -------------------------------------------------------------------------
[docs]
def manifest_to_django_orm(self) -> dict[str, Any]:
"""The VectorstoreMeta fields of the manifest."""
if not self.manifest:
raise SAMBrokerErrorNotReady(f"Manifest not loaded for {self.kind} broker.", thing=self.kind)
spec = self.manifest.spec
retval = {**super().manifest_to_django_orm()}
retval.update(
{
"spec": spec.model_dump(mode="json"),
"backend": spec.backend,
"hosting": spec.hosting,
"is_active": spec.isActive,
"dimension": spec.index.dimension,
"metric": spec.index.metric,
"deletion_protection": spec.index.deletionProtection,
"embeddings_model": spec.embeddings.model,
}
)
return retval
[docs]
def spec_of(self, vectorstore: VectorstoreMeta) -> SAMVectorstoreSpec:
"""The Vectorstore's spec, as it was applied."""
data = dict(vectorstore.spec or {})
if not data:
data = {
"backend": vectorstore.backend,
"hosting": vectorstore.hosting,
"connection": vectorstore.connection.name if vectorstore.connection else None,
"index": {"dimension": vectorstore.dimension, "metric": vectorstore.metric},
"embeddings": {
"provider": vectorstore.embeddings_provider.name if vectorstore.embeddings_provider else "",
"model": vectorstore.embeddings_model,
},
}
return SAMVectorstoreSpec(**data)
[docs]
def django_orm_to_manifest_dict(self) -> Optional[dict]:
"""The Vectorstore as a manifest, with its status."""
vectorstore = self.vectorstore
if not vectorstore:
return None
meta = SAMVectorstoreMetadata(
name=vectorstore.name,
description=vectorstore.description,
version=vectorstore.version,
tags=vectorstore.tags_list,
annotations=vectorstore.annotations if isinstance(vectorstore.annotations, list) else [],
)
status = SAMVectorstoreStatus(
accountNumber=vectorstore.user_profile.account.account_number,
username=vectorstore.user_profile.user.username,
recordLocator=vectorstore.record_locator,
created=vectorstore.created_at,
modified=vectorstore.updated_at,
vectorstoreStatus=vectorstore.status,
message=vectorstore.status_message or None,
indexName=vectorstore.index_name or None,
endpoint=vectorstore.endpoint_url or None,
apiKeySecret=vectorstore.api_key_secret.name if vectorstore.api_key_secret else None,
vectorCount=vectorstore.vector_count,
documentCount=vectorstore.documents.count(), # type: ignore[attr-defined]
snapshotCount=vectorstore.snapshots.count(), # type: ignore[attr-defined]
deployedAt=vectorstore.deployed_at,
lastCheckedAt=vectorstore.last_checked_at,
lastSnapshotAt=vectorstore.last_snapshot_at,
lastMaintenanceAt=vectorstore.last_maintenance_at,
)
model = SAMVectorstore(
apiVersion=self.api_version, kind=self.kind, metadata=meta, spec=self.spec_of(vectorstore), status=status
)
return model.model_dump(mode="json")
###########################################################################
# Smarter abstract property implementations
###########################################################################
@property
def formatted_class_name(self) -> str:
return self.formatted_text(f"{SAMVectorstoreBroker.__name__}[{id(self)}]")
@property
def kind(self) -> str:
return MANIFEST_KIND
@property
def manifest(self) -> Optional[SAMVectorstore]:
"""The Vectorstore manifest, as a Pydantic model, from the manifest loader."""
if self._manifest:
if not isinstance(self._manifest, SAMVectorstore):
raise SAMVectorstoreBrokerError("Cached manifest is not a SAMVectorstore instance", thing=self.kind)
return self._manifest
if self.loader and self.loader.manifest_kind == self.kind:
self._manifest = SAMVectorstore(
apiVersion=self.loader.manifest_api_version,
kind=self.loader.manifest_kind,
metadata=SAMVectorstoreMetadata(**self.loader.manifest_metadata),
spec=SAMVectorstoreSpec(**self.loader.manifest_spec),
)
return self._manifest
@property
def ORMMetaModelClass(self) -> Type[VectorstoreMeta]:
return VectorstoreMeta
@property
def ORMModelClass(self) -> Type[VectorstoreMeta]:
return VectorstoreMeta
@property
def orm_meta_instance(self) -> Optional[VectorstoreMeta]: # type: ignore[override]
return self.vectorstore
@property
def orm_instance(self) -> Optional[VectorstoreMeta]: # type: ignore[override]
return self.vectorstore
[docs]
def cache_invalidations(self) -> None:
if self.user_profile:
invalidate_all_cached_vectorstores_for_user_profile(self.user_profile)
if self._vectorstore:
VectorstoreMeta.get_cached_object(pk=self._vectorstore.pk, invalidate=True)
return super().cache_invalidations()
###########################################################################
# Smarter manifest abstract method implementations
###########################################################################
[docs]
def example_manifest(self, request: HttpRequest, *args, **kwargs) -> SmarterJournaledJsonResponse:
"""An example Vectorstore manifest: a self-hosted Qdrant database for a knowledge base."""
command = SmarterJournalCliCommands(self.example_manifest.__name__)
model = SAMVectorstore(
apiVersion=self.api_version,
kind=self.kind,
metadata=SAMVectorstoreMetadata(
name="example_knowledge_base",
description="A self-hosted Qdrant vector database of a company's product documentation.",
version="1.0.0",
tags=["example", "rag", "qdrant"],
annotations=[{"smarter.sh/vectorstore/purpose": "example"}],
),
spec=SAMVectorstoreSpec(
backend="qdrant",
hosting="self_hosted",
index=SAMVectorstoreIndex(dimension=1536, metric="cosine"),
embeddings=SAMVectorstoreEmbeddings(provider="openai", model="text-embedding-3-small"),
selfHosted=SAMVectorstoreSelfHosted(storage="10Gi"),
maintenance=SAMVectorstoreMaintenance(snapshotIntervalHours=24, snapshotRetention=7),
),
status=SAMVectorstoreStatus(
accountNumber=smarter_cached_objects.smarter_account.account_number,
username=smarter_cached_objects.smarter_admin.username,
recordLocator="vectorstoremeta-abc123",
created=datetime.datetime.now(),
modified=datetime.datetime.now(),
vectorstoreStatus=VectorstoreStatus.PENDING.value,
),
)
return self.json_response_ok(command=command, data=model.model_dump(mode="json"))
[docs]
def get(self, request: HttpRequest, *args, **kwargs) -> SmarterJournaledJsonResponse:
"""The Vectorstores that the user may read, optionally filtered by name."""
command = SmarterJournalCliCommands(self.get.__name__)
name = self.clean_cli_param(
param=kwargs.get(SAMMetadataKeys.NAME.value, None),
param_name="name",
url=self.smarter_build_absolute_uri(request),
)
if self.user_profile is None:
raise SAMBrokerErrorNotReady("user_profile is not set.")
vectorstores = VectorstoreMeta.objects.with_read_permission_for(self.user_profile.user) # type: ignore[attr-defined]
if name:
vectorstores = vectorstores.filter(name=name)
items = [
self.to_camel_case(VectorstoreSerializer(vectorstore).data)
for vectorstore in vectorstores.order_by("name")[:MAX_RESULTS]
]
data = {
SAMKeys.APIVERSION.value: self.api_version,
SAMKeys.KIND.value: self.kind,
SAMMetadataKeys.NAME.value: name,
SAMKeys.METADATA.value: {"count": len(items)},
SCLIResponseGet.KWARGS.value: kwargs,
SCLIResponseGet.DATA.value: {
SCLIResponseGetData.TITLES.value: self.get_model_titles(serializer=VectorstoreSerializer()),
SCLIResponseGetData.ITEMS.value: items,
},
}
return self.json_response_ok(command=command, data=data)
[docs]
def apply( # pylint: disable=too-many-locals
self, request: HttpRequest, *args, **kwargs
) -> SmarterJournaledJsonResponse:
"""
Create or update the Vectorstore.
It does not create its database: deploy does.
Once it is deployed, its backend, hosting, index, and self-hosted storage cannot change,
because the database would no longer match. If its embeddings model changes, its loaded
documents are loaded again.
"""
command = SmarterJournalCliCommands(self.apply.__name__)
if not self.ready or not self.manifest:
raise SAMBrokerErrorNotReady(
f"{self.kind} {self.name} broker is not ready", thing=self.kind, command=command
)
if not self.user_profile:
raise SAMBrokerErrorNotReady("user_profile is not set.", thing=self.kind, command=command)
if not (self.user_profile.user.is_staff or self.user_profile.user.is_superuser):
raise SAMVectorstoreBrokerError(
f"Only account admins may apply a {self.kind}.", thing=self.kind, command=command
)
spec = self.manifest.spec
provider = self.resolve_provider(spec.embeddings.provider)
connection = self.resolve_connection(spec.connection)
data = self.manifest_to_django_orm()
tags = data.pop("tags", None) or []
for field in ("id", "created_at", "updated_at"):
data.pop(field, None)
existing = VectorstoreMeta.objects.filter(
user_profile=self.user_profile, name=self.manifest.metadata.name
).first()
reload_documents = False
if existing and existing.status != VectorstoreStatus.PENDING:
previous = self.spec_of(existing)
changed = [field for field, value in IMMUTABLE_WHILE_DEPLOYED.items() if value(previous) != value(spec)]
if changed:
raise SAMVectorstoreBrokerError(
f"{self.kind} {existing.name} is deployed, so {', '.join(changed)} cannot change. "
"Delete it, and apply it again, to change them.",
thing=self.kind,
command=command,
)
reload_documents = (
previous.embeddings.provider,
previous.embeddings.model,
previous.embeddings.chunkSize,
) != (
spec.embeddings.provider,
spec.embeddings.model,
spec.embeddings.chunkSize,
)
with transaction.atomic():
vectorstore = existing or VectorstoreMeta(user_profile=self.user_profile)
for key, value in data.items():
setattr(vectorstore, key, value)
vectorstore.embeddings_provider = provider
vectorstore.connection = connection
if not vectorstore.index_name or vectorstore.status == VectorstoreStatus.PENDING:
vectorstore.index_name = spec.index.name or vectorstore.default_index_name()
try:
vectorstore.save()
vectorstore.tags.set(tags)
except Exception as e:
raise SAMVectorstoreBrokerError(
f"Failed to apply {self.kind} {self.manifest.metadata.name}: {e}", thing=self.kind, command=command
) from e
self._vectorstore = vectorstore
if (
existing
and existing.is_self_hosted
and existing.status in (VectorstoreStatus.PROVISIONING, VectorstoreStatus.READY)
):
# e.g. its cpu or memory changed.
VectorstoreService(vectorstore).backend.provision()
if reload_documents:
self.reload_documents(vectorstore)
self.cache_invalidations()
return self.json_response_ok(command=command, data=self.to_json())
[docs]
@staticmethod
def reload_documents(vectorstore: VectorstoreMeta) -> int:
"""Load the loaded documents again, e.g. with a new embeddings model."""
# pylint: disable=import-outside-toplevel
from smarter.apps.vectorstore.tasks import load_vectorstore_document
documents = list(vectorstore.documents.filter(status=VectorstoreDocumentStatus.LOADED)) # type: ignore[attr-defined]
for document in documents:
document.status = VectorstoreDocumentStatus.PENDING
document.save(update_fields=["status", "updated_at"])
if vectorstore.status == VectorstoreStatus.READY:
load_vectorstore_document.delay(document.pk)
return len(documents)
[docs]
def prompt(self, request: HttpRequest, *args, **kwargs) -> SmarterJournaledJsonResponse:
command = SmarterJournalCliCommands(self.prompt.__name__)
raise SAMBrokerErrorNotImplemented(message="Prompt not implemented", thing=self.kind, command=command)
[docs]
def describe(self, request: HttpRequest, *args, **kwargs) -> SmarterJournaledJsonResponse:
"""The Vectorstore as a manifest, with its status."""
command = SmarterJournalCliCommands(self.describe.__name__)
if self.name is None:
raise SAMBrokerErrorNotReady(f"{self.kind} name property is not set.", thing=self.kind, command=command)
if not self.vectorstore:
raise SAMBrokerErrorNotFound(f"{self.kind} {self.name} not found", thing=self.kind, command=command)
try:
data = self.django_orm_to_manifest_dict()
except Exception as e:
raise SAMVectorstoreBrokerError(
f"Failed to describe {self.kind} {self.name}: {e}", thing=self.kind, command=command
) from e
return self.json_response_ok(command=command, data=data)
[docs]
def delete(self, request: HttpRequest, *args, **kwargs) -> SmarterJournaledJsonResponse:
"""Destroy the database and its data, unless deletionProtection is enabled, then delete the Vectorstore."""
command = SmarterJournalCliCommands(self.delete.__name__)
vectorstore = self.owned_vectorstore(command)
service = VectorstoreService(vectorstore)
try:
if vectorstore.deployed_at:
service.destroy()
elif vectorstore.deletion_protection:
service.destroy() # raises
secret = vectorstore.api_key_secret
vectorstore.delete()
if secret is not None:
secret.delete()
except Exception as e:
raise SAMVectorstoreBrokerError(
f"Failed to delete {self.kind} {self.name}: {e}", thing=self.kind, command=command
) from e
self._vectorstore = None
self.cache_invalidations()
return self.json_response_ok(command=command, data={})
[docs]
def deploy(self, request: HttpRequest, *args, **kwargs) -> SmarterJournaledJsonResponse:
"""Create the database.
A self-hosted server takes a minute or two to become ready: follow it with describe.
"""
command = SmarterJournalCliCommands(self.deploy.__name__)
vectorstore = self.owned_vectorstore(command)
try:
VectorstoreService(vectorstore).deploy()
except Exception as e:
raise SAMVectorstoreBrokerError(
f"Failed to deploy {self.kind} {self.name}: {e}", thing=self.kind, command=command
) from e
self.cache_invalidations()
return self.json_response_ok(command=command, data=self.django_orm_to_manifest_dict() or {})
[docs]
def undeploy(self, request: HttpRequest, *args, **kwargs) -> SmarterJournaledJsonResponse:
"""Stop serving.
A self-hosted server's data is kept, and a managed index is left as it is.
"""
command = SmarterJournalCliCommands(self.undeploy.__name__)
vectorstore = self.owned_vectorstore(command)
try:
VectorstoreService(vectorstore).undeploy()
except Exception as e:
raise SAMVectorstoreBrokerError(
f"Failed to undeploy {self.kind} {self.name}: {e}", thing=self.kind, command=command
) from e
self.cache_invalidations()
return self.json_response_ok(command=command, data=self.django_orm_to_manifest_dict() or {})
[docs]
def logs(self, request: HttpRequest, *args, **kwargs) -> SmarterJournaledJsonResponse:
"""A self-hosted Qdrant server's recent logs."""
command = SmarterJournalCliCommands(self.logs.__name__)
vectorstore = self.owned_vectorstore(command)
logs = VectorstoreService(vectorstore).backend.logs() if vectorstore.is_self_hosted else None
return self.json_response_ok(command=command, data={"logs": logs or ""})
__all__ = ["SAMVectorstoreBroker", "SAMVectorstoreBrokerError"]