a8c0bb42a7
把 3a 的 Space 作用域推到 KB 层:KB/Doc/DocLink owner→space_id,作用域键 owner/name→space_id/name,同空间成员共享知识库(检索/入库/文库/双链/图谱)。 后端: - store: KB/Doc/DocLink 加 space_id,唯一索引 (owner,*)→(space_id,*),owner 降级创建人; 查询全改 space 作用域;SaveDoc/ListVault/GetDocByID/DeleteDocByID/ReplaceDocLinks/ ResolveInboundLinks/ListLinks 改 space;tenantIDForSpace 补异步入库租户 - MigrateKBSpaces 启动迁移(space_id 回填 + 唯一索引换新,同 Agent 顺序坑规避) - scopedKB owner/name→space_id/name;IngestJob 契约加 SpaceID;enqueueIngest/runIngest 穿 space;MinIO 对象键改 space/kb/doc(避免跨空间同名撞键,老键不透明不迁) - KB 写路由(create/ingest/ingest_file/note/delete)挂 RequireSpaceRole(member):viewer 只读 - 存储层重灌迁移端点 POST /admin/migrate-kb-storage(异步):为存量文档入队新 space 作用域的重灌作业(复用 JetStream 入库 worker 池),先删旧键;避免同步重嵌撑爆 HTTP 超时 桌面端: - KbView 收 spaceId(变则重拉库)+spaceReadOnly(viewer 禁建库/入库/文件/笔记);VaultPanel 同 验证(gateway+mcp-go+Milvus/Neo4j/embedding 全栈): - PG 迁移: 20/21 KB + 50/54 doc 回填 space_id(4 未迁=pre-多租户 owner='wt' 空租户遗留, 正确跳过),唯一索引 idx_kb_sn/idx_doc_skn 换新、旧索引删除 - KB 共享: member 见共享库 / viewer 建库·入库 403 / 切回个人空间隔离(看不到) - 全向量链路: RagA 入库(真 dashscope embedding)→ RagB(空间member)检索命中 RagA 内容 - 存储重灌: 端点异步入队 49 作业(worker 池背压处理),重灌后老文档在新 space 键可检索 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
42 lines
2.3 KiB
Go
42 lines
2.3 KiB
Go
package contract
|
|
|
|
import "encoding/json"
|
|
|
|
// 入库工作队列约定(准生产级)。把"解析→切块→向量化→图谱"从网关的裸 goroutine 升级为
|
|
// JetStream 持久工作队列:崩溃可重投、有界并发背压、削峰不 OOM。与任务流(SUNDYNIX_TASKS)同级。
|
|
const (
|
|
StreamIngest = "SUNDYNIX_INGEST" // 入库作业 JetStream 流(持久,作业不因网关重启而丢)
|
|
SubjectIngest = "sundynix.ingest" // 入库作业发布前缀;实际 sundynix.ingest.<job_id>
|
|
SubjectIngestAll = "sundynix.ingest.>" // 流捕获的通配
|
|
ConsumerIngest = "ingest-workers" // 入库 worker 持久消费者(队列组:多副本负载均衡 + 背压)
|
|
)
|
|
|
|
// IngestSubject 返回某入库作业的发布主题。
|
|
func IngestSubject(jobID string) string { return SubjectIngest + "." + jobID }
|
|
|
|
// IngestJob 是一条入库作业(claim-check 模式):大文件原始字节/正文先暂存到对象存储,
|
|
// 作业消息只带「暂存键 + 元数据」这条小消息进队列,worker 消费时再按 StageKey 取回原料。
|
|
// 这样作业消息恒小(不受 JetStream max_payload 限制),且崩溃重投只搬一个引用。
|
|
type IngestJob struct {
|
|
JobID string `json:"job_id"`
|
|
Owner string `json:"owner"` // 雪花 user.id(上传者/创建人,归属)
|
|
SpaceID string `json:"space_id"` // 所属工作区(增量3:资源共享作用域)
|
|
KBName string `json:"kb_name"` // 知识库展示名(原文留存分区)
|
|
Scoped string `json:"scoped"` // space_id/kb 作向量/全文/图谱分区键
|
|
ForceDoc string `json:"force_doc,omitempty"` // 非空=强制文档名(笔记编辑保持身份稳定)
|
|
Filename string `json:"filename,omitempty"` // 非空=文件入库(StageKey 指向原始字节,需先解析)
|
|
StageKey string `json:"stage_key"` // 暂存对象键:原始文件字节 或 纯文本
|
|
IsText bool `json:"is_text"` // true=StageKey 指向纯文本(跳过解析)
|
|
}
|
|
|
|
// Marshal/Unmarshal 入库作业(JetStream 消息体)。
|
|
func (j *IngestJob) Marshal() ([]byte, error) { return json.Marshal(j) }
|
|
|
|
func UnmarshalIngestJob(data []byte) (*IngestJob, error) {
|
|
var j IngestJob
|
|
if err := json.Unmarshal(data, &j); err != nil {
|
|
return nil, err
|
|
}
|
|
return &j, nil
|
|
}
|