diff --git a/godbrain_core/cpp_memory_store/README.md b/godbrain_core/cpp_memory_store/README.md index 33f2ed9..7805524 100644 --- a/godbrain_core/cpp_memory_store/README.md +++ b/godbrain_core/cpp_memory_store/README.md @@ -1,4 +1,4 @@ -# C++ Alexandria Memory Store (cut 10) +# C++ Alexandria Memory Store (cut 11) Same stdin JSON door as Go `memory-store.exe`. Go under `godbrain_core/memory_store/` stays. This tree is the C++ replacement; Start/Heal still launch Go until this @@ -49,6 +49,10 @@ Cut 10: `rag-eval.exe` offline hybrid fixture (`-corpus` path, default Go testdata). `-live` still uses Go. Thresholds match Go (Recall/MRR/nDCG ≥ 0.90, citation 1.0, no leakage). +Cut 11: ingest `document` + `chunks` (Local-Document-Adapter). Same gates as Go: +metadata and chunks together, `content_sha256` of `raw_transcript`, contiguous +UTF-8 byte ranges, forbidden-secret scan, immutable `chunks` collection. + Start/Heal still launch Go. ```powershell diff --git a/godbrain_core/cpp_memory_store/include/godbrain/memory_store/embedding.hpp b/godbrain_core/cpp_memory_store/include/godbrain/memory_store/embedding.hpp index 6aeb371..aa38e1e 100644 --- a/godbrain_core/cpp_memory_store/include/godbrain/memory_store/embedding.hpp +++ b/godbrain_core/cpp_memory_store/include/godbrain/memory_store/embedding.hpp @@ -48,6 +48,7 @@ bool embedding_cosine(const std::vector& left, const std::vector& bool embedding_embed( const EmbeddingRuntime& runtime, const std::string& content, std::vector* vector, std::string* err); bool embedding_embed_fake(int dimension, const std::string& content, std::vector* vector, std::string* err); +bool sha256_hex(const std::string& bytes, std::string* hex, std::string* err); int run_embedding_self_test(); diff --git a/godbrain_core/cpp_memory_store/include/godbrain/memory_store/protocol.hpp b/godbrain_core/cpp_memory_store/include/godbrain/memory_store/protocol.hpp index 6bbc513..5678b32 100644 --- a/godbrain_core/cpp_memory_store/include/godbrain/memory_store/protocol.hpp +++ b/godbrain_core/cpp_memory_store/include/godbrain/memory_store/protocol.hpp @@ -52,6 +52,30 @@ struct Claim { std::vector evidence_spans; }; +struct DocumentMetadata { + std::string source_label; + std::string display_name; + std::string file_sha256; + std::string content_sha256; + std::string extraction_method; + std::vector languages; + std::string backend; + std::string backend_version; + int chunk_count = 0; + bool has_ocr_confidence = false; + double ocr_confidence = 0; +}; + +struct SourceChunk { + int index = 0; + int count = 0; + int start_byte = 0; + int end_byte = 0; + std::string text; + bool has_confidence = false; + double confidence = 0; +}; + struct SkillExtracted { std::string name; std::string content; @@ -77,8 +101,12 @@ struct DistillationPayload { std::vector opsec_candidates; std::vector skills_extracted; bool has_document = false; + DocumentMetadata document; + std::vector chunks; }; +bool validate_document_payload(const DistillationPayload& p, std::string* err); + std::string json_escape(const std::string& s); std::string skill_procedure_text(const SkillExtracted& skill); std::string skill_stable_id(const std::string& name, const std::string& content); diff --git a/godbrain_core/cpp_memory_store/src/embedding.cpp b/godbrain_core/cpp_memory_store/src/embedding.cpp index 57394e2..b423bbb 100644 --- a/godbrain_core/cpp_memory_store/src/embedding.cpp +++ b/godbrain_core/cpp_memory_store/src/embedding.cpp @@ -146,6 +146,8 @@ std::wstring fields_join(const std::wstring& in) { return out; } +} // namespace + bool sha256_hex(const std::string& bytes, std::string* hex, std::string* err) { BCRYPT_ALG_HANDLE alg = nullptr; NTSTATUS st = BCryptOpenAlgorithmProvider(&alg, BCRYPT_SHA256_ALGORITHM, nullptr, 0); @@ -174,6 +176,8 @@ bool sha256_hex(const std::string& bytes, std::string* hex, std::string* err) { return true; } +namespace { + std::string wide_lower(const std::wstring& w) { if (w.empty()) return ""; std::wstring copy = w; diff --git a/godbrain_core/cpp_memory_store/src/protocol.cpp b/godbrain_core/cpp_memory_store/src/protocol.cpp index b3ff044..c6aa712 100644 --- a/godbrain_core/cpp_memory_store/src/protocol.cpp +++ b/godbrain_core/cpp_memory_store/src/protocol.cpp @@ -1,4 +1,5 @@ #include "godbrain/memory_store/protocol.hpp" +#include "godbrain/memory_store/embedding.hpp" #include "godbrain/memory_store/json.hpp" #include "godbrain/memory_store/state_machine.hpp" #include "godbrain/memory_store/embedding.hpp" @@ -14,6 +15,7 @@ #include #include +#include #include #include #include @@ -80,6 +82,30 @@ const char* kPayloadKeys[] = { nullptr, }; +const char* kDocumentKeys[] = { + "source_label", + "display_name", + "file_sha256", + "content_sha256", + "extraction_method", + "languages", + "backend", + "backend_version", + "chunk_count", + "ocr_confidence", + nullptr, +}; + +const char* kChunkKeys[] = { + "index", + "count", + "start_byte", + "end_byte", + "text", + "confidence", + nullptr, +}; + const char* kProvKeys[] = { "source_id", "source_type", @@ -541,11 +567,170 @@ bool validate_query_skills(const QuerySkillsRequest& r, std::string* err) { return true; } -bool validate_pre_ingestion(const DistillationPayload& p, std::string* err) { - if (p.has_document) { - if (err) *err = "document payload not in cpp memory-store protocol cut 1"; +bool is_lower_hex(const std::string& value, int byte_len) { + if (static_cast(value.size()) != byte_len * 2) return false; + for (unsigned char c : value) { + if (!((c >= '0' && c <= '9') || (c >= 'a' && c <= 'f'))) return false; + } + return true; +} + +bool is_safe_display_name(const std::string& value) { + if (value.empty() || value.size() > 128 || value == "." || value == "..") return false; + for (unsigned char c : value) { + if (c == '/' || c == '\\' || c == ':' || c <= 0x1f || (c >= 0x7f && c <= 0x9f)) return false; + } + return true; +} + +bool contains_forbidden_document_content(std::string value) { + for (char& c : value) { + if (c == '\r' || c == '\v' || c == '\f') c = '\n'; + } + std::string lower = value; + for (char& c : lower) c = static_cast(std::tolower(static_cast(c))); + if (lower.find("otpauth://") != std::string::npos || lower.find("otpauth-migration://") != std::string::npos) { + return true; + } + if (value.find("-----BEGIN ") != std::string::npos && value.find("PRIVATE KEY-----") != std::string::npos) { + return true; + } + static const std::regex kForbidden[] = { + std::regex(R"(\bbearer[ \t]+[A-Za-z0-9._~+/=-]{16,})", std::regex::icase), + std::regex(R"(\b(?:AKIA|ASIA)[A-Z0-9]{16}\b)"), + std::regex(R"(\b(?:gh[pousr]_[A-Za-z0-9]{36,255}|github_pat_[A-Za-z0-9_]{40,255})\b)"), + std::regex(R"(\bAIza[0-9A-Za-z_-]{35}\b)"), + std::regex(R"(\bxox[baprs]-[A-Za-z0-9-]{20,}\b)"), + std::regex(R"(\b(?:sk_live_|rk_live_)[0-9A-Za-z]{16,}\b)"), + std::regex(R"(\beyJ[0-9A-Za-z_-]{8,}\.[0-9A-Za-z_-]{8,}\.[0-9A-Za-z_-]{8,}\b)"), + std::regex( + R"((?:api[_-]?key|password|passwd|private[_-]?key|client[_-]?secret|access[_-]?token|bearer[_-]?token)[\s]*[:=][\s]*["']?[^\s"']{8,})", + std::regex::icase), + }; + for (const auto& re : kForbidden) { + if (std::regex_search(value, re)) return true; + } + return false; +} + +bool is_local_document(const DistillationPayload& p) { + std::string extractor = p.extractor_id.empty() ? kDefaultExtractorID : p.extractor_id; + return extractor == "Local-Document-Adapter" || p.provenance.source_type == "local_document" || + p.provenance.source_id.rfind("local-document:", 0) == 0 || p.has_document || !p.chunks.empty(); +} + +bool validate_document_payload(const DistillationPayload& p, std::string* err) { + if (!is_local_document(p)) return true; + if (!p.has_document || p.chunks.empty()) { + if (err) *err = "document metadata and chunks must be provided together"; return false; } + const DocumentMetadata& d = p.document; + std::string extractor = p.extractor_id.empty() ? kDefaultExtractorID : p.extractor_id; + if (!safe_extractor_id(extractor) || p.extractor_version.empty() || p.extractor_version.size() > 128 || + p.schema_version.empty() || p.schema_version.size() > 64) { + if (err) *err = "document extractor identity is invalid"; + return false; + } + if (d.source_label.empty() || d.display_name.empty() || d.extraction_method.empty() || d.backend.empty() || + d.backend_version.empty() || d.languages.empty()) { + if (err) *err = "document provenance fields must not be blank"; + return false; + } + if (d.source_label.size() > 64 || d.display_name.size() > 128 || d.extraction_method.size() > 64 || + d.backend.size() > 64 || d.backend_version.size() > 64 || d.languages.size() > 8) { + if (err) *err = "document provenance exceeds field bounds"; + return false; + } + if (!safe_extractor_id(d.source_label) || !is_safe_display_name(d.display_name)) { + if (err) *err = "document source label or display name is unsafe"; + return false; + } + if (!is_lower_hex(d.file_sha256, 32) || !is_lower_hex(d.content_sha256, 32)) { + if (err) *err = "document SHA-256 fields must be lowercase hexadecimal"; + return false; + } + if (d.has_ocr_confidence && (d.ocr_confidence < 0 || d.ocr_confidence > 1)) { + if (err) *err = "document OCR confidence is out of bounds"; + return false; + } + const std::string expected_id = "local-document:" + d.source_label + ":" + d.display_name; + if (p.provenance.source_id != expected_id) { + if (err) *err = "document metadata does not match its safe source identity"; + return false; + } + std::string langs; + for (size_t i = 0; i < d.languages.size(); ++i) { + if (i) langs.push_back(','); + langs += d.languages[i]; + } + if (p.provenance.source_type != "local_document" || p.provenance.language != langs) { + if (err) *err = "document provenance does not match its safe source identity"; + return false; + } + std::string digest; + std::string herr; + if (!sha256_hex(p.raw_transcript, &digest, &herr) || digest != d.content_sha256) { + if (err) *err = "document content_sha256 does not match raw_transcript"; + return false; + } + if (d.chunk_count != static_cast(p.chunks.size()) || p.chunks.size() > 256) { + if (err) *err = "document chunk_count does not match bounded chunks"; + return false; + } + if (contains_forbidden_document_content(p.raw_transcript) || + contains_forbidden_document_content(d.source_label) || + contains_forbidden_document_content(d.display_name)) { + if (err) *err = "document contains forbidden sensitive content"; + return false; + } + int previous_end = 0; + const std::string& raw = p.raw_transcript; + for (size_t i = 0; i < p.chunks.size(); ++i) { + const SourceChunk& ch = p.chunks[i]; + if (ch.index != static_cast(i) || ch.count != static_cast(p.chunks.size())) { + if (err) *err = "document chunks must have contiguous indexes and a consistent count"; + return false; + } + if (ch.start_byte < 0 || ch.end_byte <= ch.start_byte || ch.end_byte > static_cast(raw.size()) || + ch.end_byte - ch.start_byte > 32 * 1024) { + if (err) *err = "document chunk byte range is invalid"; + return false; + } + if (raw.substr(static_cast(ch.start_byte), static_cast(ch.end_byte - ch.start_byte)) != + ch.text) { + if (err) *err = "document chunk text does not match raw_transcript byte range"; + return false; + } + std::string gap = raw.substr(static_cast(previous_end), static_cast(ch.start_byte - previous_end)); + bool gap_ws = true; + for (unsigned char c : gap) { + if (std::isspace(c) == 0) { + gap_ws = false; + break; + } + } + if (ch.start_byte < previous_end || !gap_ws) { + if (err) *err = "document chunks overlap or omit non-whitespace content"; + return false; + } + if (ch.has_confidence && (ch.confidence < 0 || ch.confidence > 1)) { + if (err) *err = "document chunk confidence is out of bounds"; + return false; + } + previous_end = ch.end_byte; + } + std::string trail = raw.substr(static_cast(previous_end)); + for (unsigned char c : trail) { + if (std::isspace(c) == 0) { + if (err) *err = "document chunks omit trailing non-whitespace content"; + return false; + } + } + return true; +} + +bool validate_pre_ingestion(const DistillationPayload& p, std::string* err) { std::string extractor = p.extractor_id.empty() ? kDefaultExtractorID : p.extractor_id; if (!safe_extractor_id(extractor) || p.extractor_version.empty() || p.extractor_version.size() > 128 || p.schema_version.empty() || @@ -569,7 +754,7 @@ bool validate_pre_ingestion(const DistillationPayload& p, std::string* err) { } return false; } - return true; + return validate_document_payload(p, err); } bool classify_and_parse(const std::string& json_text, Route* route, std::string* err) { @@ -843,7 +1028,70 @@ bool classify_and_parse(const std::string& json_text, Route* route, std::string* } } - route->ingest.has_document = json_has(root, "document") || json_has(root, "chunks"); + if (json_has(root, "document")) { + const Json* doc = json_get(root, "document"); + if (doc == nullptr || !json_is_object(*doc)) { + if (err) *err = "document metadata and chunks must be provided together"; + return false; + } + if (!json_reject_unknown_keys(*doc, kDocumentKeys, err)) return false; + json_string(*doc, "source_label", &route->ingest.document.source_label); + json_string(*doc, "display_name", &route->ingest.document.display_name); + json_string(*doc, "file_sha256", &route->ingest.document.file_sha256); + json_string(*doc, "content_sha256", &route->ingest.document.content_sha256); + json_string(*doc, "extraction_method", &route->ingest.document.extraction_method); + json_string(*doc, "backend", &route->ingest.document.backend); + json_string(*doc, "backend_version", &route->ingest.document.backend_version); + double cc = 0; + if (json_has(*doc, "chunk_count") && json_number(*doc, "chunk_count", &cc)) { + route->ingest.document.chunk_count = static_cast(cc); + } + if (json_has(*doc, "ocr_confidence")) { + route->ingest.document.has_ocr_confidence = true; + json_number(*doc, "ocr_confidence", &route->ingest.document.ocr_confidence); + } + const Json* langs = json_get(*doc, "languages"); + if (langs != nullptr) { + if (langs->kind != Json::Kind::Array) { + if (err) *err = "document provenance fields must not be blank"; + return false; + } + for (const Json& v : langs->arr) { + if (v.kind != Json::Kind::String) { + if (err) *err = "document provenance fields must not be blank"; + return false; + } + route->ingest.document.languages.push_back(v.str); + } + } + route->ingest.has_document = true; + } + if (json_has(root, "chunks")) { + const Json* chunks = json_get(root, "chunks"); + if (chunks == nullptr || chunks->kind != Json::Kind::Array) { + if (err) *err = "document metadata and chunks must be provided together"; + return false; + } + for (const Json& item : chunks->arr) { + if (!json_is_object(item)) { + if (err) *err = "document chunks must have contiguous indexes and a consistent count"; + return false; + } + if (!json_reject_unknown_keys(item, kChunkKeys, err)) return false; + SourceChunk ch; + double n = 0; + if (json_number(item, "index", &n)) ch.index = static_cast(n); + if (json_number(item, "count", &n)) ch.count = static_cast(n); + if (json_number(item, "start_byte", &n)) ch.start_byte = static_cast(n); + if (json_number(item, "end_byte", &n)) ch.end_byte = static_cast(n); + json_string(item, "text", &ch.text); + if (json_has(item, "confidence")) { + ch.has_confidence = true; + json_number(item, "confidence", &ch.confidence); + } + route->ingest.chunks.push_back(std::move(ch)); + } + } return validate_pre_ingestion(route->ingest, err); } @@ -974,6 +1222,41 @@ int run_self_test() { err.clear(); check(!classify_and_parse(bad_skills, &r, &err), "skills-extracted-type"); + const std::string doc_text = "hello world notes."; + const std::string doc_k = keccak256_hex(doc_text); + std::string doc_sha; + check(sha256_hex(doc_text, &doc_sha, &err), "doc-sha"); + const std::string file_sha = doc_sha; + const std::string doc_json = + std::string("{\"extractor_id\":\"Local-Document-Adapter\",\"extractor_version\":\"v1\",") + + "\"schema_version\":\"1.0\",\"raw_transcript\":\"" + doc_text + + "\",\"payload\":{\"trust_tier\":\"candidate\",\"provenance\":{" + "\"source_id\":\"local-document:notes.txt:notes.txt\",\"source_type\":\"local_document\"," + "\"source_hash\":\"" + + doc_k + "\",\"language\":\"en\"}},\"document\":{\"source_label\":\"notes.txt\"," + "\"display_name\":\"notes.txt\",\"file_sha256\":\"" + + file_sha + "\",\"content_sha256\":\"" + doc_sha + + "\",\"extraction_method\":\"text\",\"languages\":[\"en\"],\"backend\":\"utf8\"," + "\"backend_version\":\"1\",\"chunk_count\":1},\"chunks\":[{\"index\":0,\"count\":1," + "\"start_byte\":0,\"end_byte\":" + + std::to_string(doc_text.size()) + ",\"text\":\"" + doc_text + "\"}]}"; + err.clear(); + check(classify_and_parse(doc_json, &r, &err) && r.ingest.chunks.size() == 1, "document-ok"); + err.clear(); + check(!classify_and_parse( + std::string("{\"extractor_id\":\"Local-Document-Adapter\",\"extractor_version\":\"v1\",") + + "\"schema_version\":\"1.0\",\"raw_transcript\":\"" + doc_text + + "\",\"payload\":{\"trust_tier\":\"candidate\",\"provenance\":{" + "\"source_id\":\"local-document:notes.txt:notes.txt\",\"source_type\":\"local_document\"," + "\"source_hash\":\"" + + doc_k + "\",\"language\":\"en\"}},\"document\":{\"source_label\":\"notes.txt\"," + "\"display_name\":\"notes.txt\",\"file_sha256\":\"" + + file_sha + "\",\"content_sha256\":\"" + doc_sha + + "\",\"extraction_method\":\"text\",\"languages\":[\"en\"],\"backend\":\"utf8\"," + "\"backend_version\":\"1\",\"chunk_count\":0}}", + &r, &err), + "document-chunks-required"); + check(allowed_run_transition(kRunStaging, kRunValidated), "run-stag-val"); check(allowed_run_transition(kRunStaging, kRunFailed), "run-stag-fail"); check(allowed_run_transition(kRunValidated, kRunCommitted), "run-val-com"); diff --git a/godbrain_core/cpp_memory_store/src/store.cpp b/godbrain_core/cpp_memory_store/src/store.cpp index 3367d9c..a1efbf0 100644 --- a/godbrain_core/cpp_memory_store/src/store.cpp +++ b/godbrain_core/cpp_memory_store/src/store.cpp @@ -380,6 +380,19 @@ static bool ensure_indexes(StoreHandle* h, std::string* err) { bson_destroy(&keys); if (!ok) return false; } + { + bson_t keys = BSON_INITIALIZER; + BSON_APPEND_INT32(&keys, "source_hash", 1); + BSON_APPEND_INT32(&keys, "extractor_id", 1); + BSON_APPEND_INT32(&keys, "extractor_version", 1); + BSON_APPEND_INT32(&keys, "schema_version", 1); + BSON_APPEND_INT32(&keys, "chunk_index", 1); + bson_t* opt = BCON_NEW("unique", BCON_BOOL(true), "name", "source_chunk_extractor_identity"); + bool ok = create_unique_index(h, "chunks", &keys, opt, err, "chunks index"); + bson_destroy(opt); + bson_destroy(&keys); + if (!ok) return false; + } return true; } @@ -553,6 +566,29 @@ bool store_ingest( if (!retry_of.empty()) { BSON_APPEND_UTF8(&set_on, "retry_of", retry_of.c_str()); } + bson_t doc_meta = BSON_INITIALIZER; + if (payload.has_document) { + BSON_APPEND_UTF8(&doc_meta, "source_label", payload.document.source_label.c_str()); + BSON_APPEND_UTF8(&doc_meta, "display_name", payload.document.display_name.c_str()); + BSON_APPEND_UTF8(&doc_meta, "file_sha256", payload.document.file_sha256.c_str()); + BSON_APPEND_UTF8(&doc_meta, "content_sha256", payload.document.content_sha256.c_str()); + BSON_APPEND_UTF8(&doc_meta, "extraction_method", payload.document.extraction_method.c_str()); + bson_t langs = BSON_INITIALIZER; + for (size_t i = 0; i < payload.document.languages.size(); ++i) { + char idx[16]; + std::snprintf(idx, sizeof idx, "%zu", i); + BSON_APPEND_UTF8(&langs, idx, payload.document.languages[i].c_str()); + } + BSON_APPEND_ARRAY(&doc_meta, "languages", &langs); + bson_destroy(&langs); + BSON_APPEND_UTF8(&doc_meta, "backend", payload.document.backend.c_str()); + BSON_APPEND_UTF8(&doc_meta, "backend_version", payload.document.backend_version.c_str()); + BSON_APPEND_INT32(&doc_meta, "chunk_count", payload.document.chunk_count); + if (payload.document.has_ocr_confidence) { + BSON_APPEND_DOUBLE(&doc_meta, "ocr_confidence", payload.document.ocr_confidence); + } + BSON_APPEND_DOCUMENT(&set_on, "document", &doc_meta); + } bson_t update = BSON_INITIALIZER; BSON_APPEND_DOCUMENT(&update, "$setOnInsert", &set_on); @@ -566,6 +602,7 @@ bool store_ingest( bson_error_t error{}; bool ok = mongoc_collection_find_and_modify_with_opts(runs, &filter, fm, &reply, &error); mongoc_find_and_modify_opts_destroy(fm); + bson_destroy(&doc_meta); bson_destroy(&set_on); bson_destroy(&update); bson_t* run_doc = nullptr; @@ -622,6 +659,9 @@ bool store_ingest( BSON_APPEND_UTF8(&oq, "extractor_id", extractor.c_str()); BSON_APPEND_UTF8(&oq, "extractor_version", payload.extractor_version.c_str()); BSON_APPEND_UTF8(&oq, "schema_version", payload.schema_version.c_str()); + if (payload.has_document) { + BSON_APPEND_UTF8(&oq, "document.file_sha256", payload.document.file_sha256.c_str()); + } bson_t oset = BSON_INITIALIZER; BSON_APPEND_UTF8(&oset, "source_hash", source_hash.c_str()); BSON_APPEND_UTF8(&oset, "external_source_id", payload.provenance.source_id.c_str()); @@ -630,12 +670,36 @@ bool store_ingest( BSON_APPEND_UTF8(&oset, "schema_version", payload.schema_version.c_str()); BSON_APPEND_UTF8(&oset, "run_id", run_id.c_str()); BSON_APPEND_DATE_TIME(&oset, "created_at", now); + bson_t obs_doc = BSON_INITIALIZER; + if (payload.has_document) { + BSON_APPEND_UTF8(&obs_doc, "source_label", payload.document.source_label.c_str()); + BSON_APPEND_UTF8(&obs_doc, "display_name", payload.document.display_name.c_str()); + BSON_APPEND_UTF8(&obs_doc, "file_sha256", payload.document.file_sha256.c_str()); + BSON_APPEND_UTF8(&obs_doc, "content_sha256", payload.document.content_sha256.c_str()); + BSON_APPEND_UTF8(&obs_doc, "extraction_method", payload.document.extraction_method.c_str()); + bson_t langs = BSON_INITIALIZER; + for (size_t i = 0; i < payload.document.languages.size(); ++i) { + char idx[16]; + std::snprintf(idx, sizeof idx, "%zu", i); + BSON_APPEND_UTF8(&langs, idx, payload.document.languages[i].c_str()); + } + BSON_APPEND_ARRAY(&obs_doc, "languages", &langs); + bson_destroy(&langs); + BSON_APPEND_UTF8(&obs_doc, "backend", payload.document.backend.c_str()); + BSON_APPEND_UTF8(&obs_doc, "backend_version", payload.document.backend_version.c_str()); + BSON_APPEND_INT32(&obs_doc, "chunk_count", payload.document.chunk_count); + if (payload.document.has_ocr_confidence) { + BSON_APPEND_DOUBLE(&obs_doc, "ocr_confidence", payload.document.ocr_confidence); + } + BSON_APPEND_DOCUMENT(&oset, "document", &obs_doc); + } bson_t oupd = BSON_INITIALIZER; BSON_APPEND_DOCUMENT(&oupd, "$setOnInsert", &oset); bson_t oopts = BSON_INITIALIZER; BSON_APPEND_BOOL(&oopts, "upsert", true); ok = mongoc_collection_update_one(obs, &oq, &oupd, &oopts, nullptr, &error); bson_destroy(&oq); + bson_destroy(&obs_doc); bson_destroy(&oset); bson_destroy(&oupd); bson_destroy(&oopts); @@ -697,6 +761,76 @@ bool store_ingest( if (err) *err = "failed to resolve immutable source"; return false; } + if (!payload.chunks.empty()) { + mongoc_collection_t* chunks = coll(h, "chunks"); + for (const SourceChunk& ch : payload.chunks) { + bson_t fq = BSON_INITIALIZER; + BSON_APPEND_UTF8(&fq, "source_hash", source_hash.c_str()); + BSON_APPEND_UTF8(&fq, "extractor_id", extractor.c_str()); + BSON_APPEND_UTF8(&fq, "extractor_version", payload.extractor_version.c_str()); + BSON_APPEND_UTF8(&fq, "schema_version", payload.schema_version.c_str()); + BSON_APPEND_INT32(&fq, "chunk_index", ch.index); + bson_t setc = BSON_INITIALIZER; + BSON_APPEND_UTF8(&setc, "source_hash", source_hash.c_str()); + BSON_APPEND_UTF8(&setc, "extractor_id", extractor.c_str()); + BSON_APPEND_UTF8(&setc, "extractor_version", payload.extractor_version.c_str()); + BSON_APPEND_UTF8(&setc, "schema_version", payload.schema_version.c_str()); + BSON_APPEND_INT32(&setc, "chunk_index", ch.index); + BSON_APPEND_INT32(&setc, "start_byte", ch.start_byte); + BSON_APPEND_INT32(&setc, "end_byte", ch.end_byte); + BSON_APPEND_UTF8(&setc, "text", ch.text.c_str()); + BSON_APPEND_INT32(&setc, "chunk_count", ch.count); + if (ch.has_confidence) BSON_APPEND_DOUBLE(&setc, "confidence", ch.confidence); + bson_t cupd = BSON_INITIALIZER; + BSON_APPEND_DOCUMENT(&cupd, "$setOnInsert", &setc); + bson_t copts = BSON_INITIALIZER; + BSON_APPEND_BOOL(&copts, "upsert", true); + bool cok = mongoc_collection_update_one(chunks, &fq, &cupd, &copts, nullptr, &error); + bson_destroy(&fq); + bson_destroy(&setc); + bson_destroy(&cupd); + bson_destroy(&copts); + if (!cok && error.code != 11000) { + mongoc_collection_destroy(chunks); + fail_created_run(h, run_id, used_lease, "stage chunks"); + return bson_ok(false, error, err, "stage source chunks"); + } + bson_t rq = BSON_INITIALIZER; + BSON_APPEND_UTF8(&rq, "source_hash", source_hash.c_str()); + BSON_APPEND_UTF8(&rq, "extractor_id", extractor.c_str()); + BSON_APPEND_UTF8(&rq, "extractor_version", payload.extractor_version.c_str()); + BSON_APPEND_UTF8(&rq, "schema_version", payload.schema_version.c_str()); + BSON_APPEND_INT32(&rq, "chunk_index", ch.index); + mongoc_cursor_t* ccur = mongoc_collection_find_with_opts(chunks, &rq, nullptr, nullptr); + const bson_t* stored = nullptr; + bool got = mongoc_cursor_next(ccur, &stored); + std::string stored_text; + int stored_start = -1, stored_end = -1, stored_count = -1; + if (got) { + iter_utf8(stored, "text", &stored_text); + bson_iter_t cit; + if (bson_iter_init_find(&cit, stored, "start_byte") && BSON_ITER_HOLDS_INT32(&cit)) { + stored_start = bson_iter_int32(&cit); + } + if (bson_iter_init_find(&cit, stored, "end_byte") && BSON_ITER_HOLDS_INT32(&cit)) { + stored_end = bson_iter_int32(&cit); + } + if (bson_iter_init_find(&cit, stored, "chunk_count") && BSON_ITER_HOLDS_INT32(&cit)) { + stored_count = bson_iter_int32(&cit); + } + } + mongoc_cursor_destroy(ccur); + bson_destroy(&rq); + if (!got || stored_text != ch.text || stored_start != ch.start_byte || stored_end != ch.end_byte || + stored_count != ch.count) { + mongoc_collection_destroy(chunks); + fail_created_run(h, run_id, used_lease, "chunk conflict"); + if (err) *err = "stored source chunk conflicts with immutable chunk content"; + return false; + } + } + mongoc_collection_destroy(chunks); + } bson_t run_q = BSON_INITIALIZER; BSON_APPEND_UTF8(&run_q, "run_id", run_id.c_str());