393 files
This commit is contained in:
@@ -14,6 +14,7 @@ serving view, sensor and frontend build on top of it.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import hashlib
|
||||
import logging
|
||||
import math
|
||||
@@ -203,6 +204,18 @@ class DocumentStore:
|
||||
# every new blob, told when the last reference goes. Optional so the
|
||||
# store stays usable on its own (tests, tooling).
|
||||
self.text_index: DocumentTextIndex | None = None
|
||||
# Serialises every blob-registry transition that spans an await (the
|
||||
# executor write/delete). Without it an upload of content whose last
|
||||
# document was being deleted at the same moment registered the blob
|
||||
# while the delete removed the file underneath it — a document
|
||||
# pointing at nothing (bug audit 2026-09-26, SEC-6).
|
||||
self._blob_lock = asyncio.Lock()
|
||||
|
||||
@property
|
||||
def blob_lock(self) -> asyncio.Lock:
|
||||
"""The blob-registry lock, for a caller that writes blobs AND
|
||||
registers them itself (the documents-archive restore)."""
|
||||
return self._blob_lock
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Paths
|
||||
@@ -331,15 +344,18 @@ class DocumentStore:
|
||||
if object_doc_count >= MAX_DOCS_PER_OBJECT:
|
||||
raise ValueError("too_many_documents")
|
||||
|
||||
digest, wrote_new = await self.hass.async_add_executor_job(self._store_blob_sync, content)
|
||||
async with self._blob_lock:
|
||||
digest, wrote_new = await self.hass.async_add_executor_job(self._store_blob_sync, content)
|
||||
|
||||
# Register / adopt the blob and bump its refcount.
|
||||
blob = self.blobs.get(digest)
|
||||
deduped = blob is not None
|
||||
if blob is None:
|
||||
blob = {"size": len(content), "mime": mime, "refcount": 0}
|
||||
self.blobs[digest] = blob
|
||||
blob["refcount"] += 1
|
||||
# Register / adopt the blob and bump its refcount — in the same
|
||||
# critical section as the write, so a concurrent last-reference
|
||||
# delete cannot remove the file between the two.
|
||||
blob = self.blobs.get(digest)
|
||||
deduped = blob is not None
|
||||
if blob is None:
|
||||
blob = {"size": len(content), "mime": mime, "refcount": 0}
|
||||
self.blobs[digest] = blob
|
||||
blob["refcount"] += 1
|
||||
|
||||
duplicate_in_object = next(
|
||||
(did for did, d in self.documents.items() if d.get("object_id") == object_id and d.get("hash") == digest),
|
||||
@@ -386,9 +402,16 @@ class DocumentStore:
|
||||
if path.exists():
|
||||
return digest, False
|
||||
self._blobs_dir.mkdir(parents=True, exist_ok=True)
|
||||
tmp = path.with_name(path.name + ".tmp")
|
||||
tmp.write_bytes(content)
|
||||
os.replace(tmp, path) # atomic
|
||||
# A temp name of its own per write: two uploads of the same content
|
||||
# shared "<digest>.tmp", and the second os.replace found it already
|
||||
# moved away — a 500 for a perfectly good upload (bug audit
|
||||
# 2026-09-26). Replacing an identical blob is harmless.
|
||||
tmp = path.with_name(f"{path.name}.{uuid4().hex}.tmp")
|
||||
try:
|
||||
tmp.write_bytes(content)
|
||||
os.replace(tmp, path) # atomic
|
||||
finally:
|
||||
tmp.unlink(missing_ok=True)
|
||||
return digest, True
|
||||
|
||||
async def async_add_weblink(
|
||||
@@ -716,24 +739,47 @@ class DocumentStore:
|
||||
return freed
|
||||
|
||||
async def _deref_blob(self, doc: dict[str, Any]) -> int:
|
||||
"""Decrement a file doc's blob refcount; delete the blob at 0."""
|
||||
"""Decrement a file doc's blob refcount; delete the blob at 0.
|
||||
|
||||
Under the blob lock (see ``_blob_lock``): the file delete must not
|
||||
interleave with an upload adopting the same content.
|
||||
"""
|
||||
if doc.get("kind") != KIND_FILE:
|
||||
return 0
|
||||
digest = doc.get("hash")
|
||||
if not isinstance(digest, str):
|
||||
return 0
|
||||
blob = self.blobs.get(digest)
|
||||
if blob is None:
|
||||
async with self._blob_lock:
|
||||
blob = self.blobs.get(digest)
|
||||
if blob is None:
|
||||
return 0
|
||||
blob["refcount"] = blob.get("refcount", 1) - 1
|
||||
if blob["refcount"] <= 0:
|
||||
size = int(blob.get("size", 0))
|
||||
self.blobs.pop(digest, None)
|
||||
if self.text_index is not None:
|
||||
await self.text_index.async_forget(digest)
|
||||
await self.hass.async_add_executor_job(self._delete_blob_sync, digest)
|
||||
return size
|
||||
return 0
|
||||
blob["refcount"] = blob.get("refcount", 1) - 1
|
||||
if blob["refcount"] <= 0:
|
||||
size = int(blob.get("size", 0))
|
||||
self.blobs.pop(digest, None)
|
||||
if self.text_index is not None:
|
||||
await self.text_index.async_forget(digest)
|
||||
await self.hass.async_add_executor_job(self._delete_blob_sync, digest)
|
||||
return size
|
||||
return 0
|
||||
|
||||
async def async_rehome(self, doc_ids: set[str], object_id: str) -> int:
|
||||
"""Re-stamp documents onto another object. Returns the docs moved.
|
||||
|
||||
A shared spare-part pool that moves to a borrower when its owner is
|
||||
deleted (#111) takes its manual along: the part keeps its ``doc_id``,
|
||||
and the document must not go down with the old owner (bug audit
|
||||
2026-09-26, SEC-10). Unknown ids are skipped.
|
||||
"""
|
||||
moved = 0
|
||||
for doc_id in doc_ids:
|
||||
doc = self.documents.get(doc_id)
|
||||
if doc is not None and doc.get("object_id") != object_id:
|
||||
doc["object_id"] = object_id
|
||||
moved += 1
|
||||
if moved:
|
||||
await self._async_save()
|
||||
return moved
|
||||
|
||||
def _delete_blob_sync(self, digest: str) -> None:
|
||||
"""Delete a blob file and its extracted-text sidecar (executor)."""
|
||||
@@ -848,27 +894,58 @@ class DocumentStore:
|
||||
"""
|
||||
issues = await self.async_find_issues()
|
||||
freed = 0
|
||||
orphans = zero = dangling = 0
|
||||
# Every deletion re-checks its finding UNDER the blob lock: the scan
|
||||
# above ran without it, and an upload of the same content (file
|
||||
# written, registry entry pending) looked like an orphan — the
|
||||
# cleanup deleted the file under the fresh document (bug audit
|
||||
# 2026-09-27; uploads and deletes already serialised on the lock).
|
||||
for digest in issues["orphan_blobs"]:
|
||||
freed += await self.hass.async_add_executor_job(self._reclaim_blob_sync, digest)
|
||||
async with self._blob_lock:
|
||||
if digest in self.blobs:
|
||||
continue # registered meanwhile — not an orphan any more
|
||||
freed += await self.hass.async_add_executor_job(self._reclaim_blob_sync, digest)
|
||||
orphans += 1
|
||||
for digest in issues["zero_refcount"]:
|
||||
blob = self.blobs.pop(digest, None)
|
||||
freed += int(blob.get("size", 0)) if blob else 0
|
||||
await self.hass.async_add_executor_job(self._delete_blob_sync, digest)
|
||||
async with self._blob_lock:
|
||||
blob = self.blobs.get(digest)
|
||||
if blob is None or int(blob.get("refcount", 0)) > 0:
|
||||
continue # adopted by an upload meanwhile
|
||||
self.blobs.pop(digest, None)
|
||||
freed += int(blob.get("size", 0))
|
||||
await self.hass.async_add_executor_job(self._delete_blob_sync, digest)
|
||||
zero += 1
|
||||
for did in issues["dangling_docs"]:
|
||||
doc = self.documents.pop(did, None)
|
||||
if doc is not None:
|
||||
# Reconcile the now over-counted blob refcount so no phantom
|
||||
# registry entry lingers; the file is already gone (0 real bytes).
|
||||
await self._deref_blob(doc)
|
||||
doc = self.documents.get(did)
|
||||
if doc is None:
|
||||
continue
|
||||
async with self._blob_lock:
|
||||
doc_hash = doc.get("hash")
|
||||
if isinstance(doc_hash, str) and doc_hash in self.blobs:
|
||||
if await self.hass.async_add_executor_job(self._blob_file_exists, doc_hash):
|
||||
continue # its content arrived meanwhile (upload / restore)
|
||||
self.documents.pop(did, None)
|
||||
# Reconcile the now over-counted blob refcount so no phantom
|
||||
# registry entry lingers; the file is already gone (0 real bytes).
|
||||
# Outside the lock: _deref_blob takes it itself.
|
||||
await self._deref_blob(doc)
|
||||
dangling += 1
|
||||
if any(issues.values()):
|
||||
await self._async_save()
|
||||
return {
|
||||
"orphans_deleted": len(issues["orphan_blobs"]),
|
||||
"zero_refcount_cleared": len(issues["zero_refcount"]),
|
||||
"dangling_removed": len(issues["dangling_docs"]),
|
||||
"orphans_deleted": orphans,
|
||||
"zero_refcount_cleared": zero,
|
||||
"dangling_removed": dangling,
|
||||
"bytes_freed": freed,
|
||||
}
|
||||
|
||||
def _blob_file_exists(self, digest: str) -> bool:
|
||||
"""Whether the blob file is on disk (executor; invalid digest = no)."""
|
||||
try:
|
||||
return self.blob_path(digest).is_file()
|
||||
except ValueError:
|
||||
return False
|
||||
|
||||
def _reclaim_blob_sync(self, digest: str) -> int:
|
||||
"""Delete an orphan blob file, returning its size (0 if already gone)."""
|
||||
try:
|
||||
|
||||
Reference in New Issue
Block a user