Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions bootstrap/cli/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,9 @@ def call(self, verb: str, payload: dict[str, Any]) -> tuple[int, dict[str, Any]]
def healthz(self) -> tuple[int, dict[str, Any]]:
return 200, {"status": "ok", "profile": self._srv.config.profile}

def close(self) -> None:
self._srv.close(wait=True)


class HttpClient:
"""Drive a running ``bootstrap`` server over HTTP (``POST /v1/<verb>``)."""
Expand Down
101 changes: 99 additions & 2 deletions bootstrap/core/handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

import os
import sys
import uuid
from datetime import datetime
from importlib import import_module
from typing import Any, Callable
Expand Down Expand Up @@ -45,6 +46,7 @@
Modality = _type_def_module.Modality
Scope = _type_def_module.Scope
EvolveMode = import_module("jiuwen_memory.construction").EvolveMode
INGEST_JOB_PREFIX = import_module("jiuwen_memory.control.ingest_job").INGEST_JOB_PREFIX

_control_types_module = import_module("jiuwen_memory.control.types")
Action = _control_types_module.Action
Expand Down Expand Up @@ -371,12 +373,103 @@ def _usage_view(usage) -> Body:
}


def _video_unit_view(unit: MemoryUnit) -> Body:
"""Serialize fields needed by video job responses."""
return {
**_unit_view(unit),
"source": unit.source.value,
"source_ref": unit.source_ref,
"provenance": list(unit.provenance),
"metadata": dict(unit.metadata),
}


def _submit_video(srv, payload: Body, *, scope: Scope, identity: Scope) -> Body:
"""Submit a video write through the Control-managed ingest queue."""
uri = str(_require(payload, "uri")).strip()
if not uri:
raise ValidationError("missing required field: 'uri'")
payload_id = str(payload.get("payload_id") or uuid.uuid4())
raw_meta = payload.get("metadata")
if not isinstance(raw_meta, dict):
raw_meta = {}
metadata = dict(raw_meta)
metadata.update(
{"infer": "true", "pipeline": "video", "payload_id": payload_id}
)
requested_assets = payload.get("assets")
extras = requested_assets if isinstance(requested_assets, list) else []
assets = [uri, *(str(item) for item in extras if str(item) != uri)]

submission = srv.ingest_jobs.submit(
payload_id=payload_id,
source_ref=uri,
scope=scope,
task=lambda: srv.api.add(
uri,
scope,
Modality.VIDEO,
identity=identity,
assets=assets,
tags=payload.get("tags"),
metadata=metadata,
),
)
job = submission.job
return {
"ok": True,
"op": "add",
"accepted": True,
"job_id": job.id,
"video_id": payload_id,
"status": job.status,
"reused": submission.reused,
"feedback_message": (
"已返回该视频现有的处理任务。"
if submission.reused
else "视频处理任务已提交。"
),
}


def _ingest_job_status(srv, job, *, identity: Scope) -> Body:
"""Adapt a Control ingest job to the shared job response shape."""
scope = job.scope
units = []
for unit_id in job.unit_ids:
try:
units.append(srv.api.get(unit_id, scope, identity=identity))
except NotFoundError:
continue
items = [_video_unit_view(unit) for unit in units]
body: Body = {
"ok": True,
"op": "job",
"job_id": job.id,
"video_id": job.payload_id,
"status": job.status,
"count": len(items),
"item_ids": [item["item_id"] for item in items],
"items": items,
}
if job.status == "succeeded":
body["feedback_message"] = f"视频处理完成,共生成 {len(items)} 条多模态记忆。"
elif job.status == "failed":
body["error"] = job.error
body["feedback_message"] = "视频处理失败。"
else:
body["feedback_message"] = "视频正在处理中。"
return body


# --- per-verb handlers ----------------------------------------------------- #


def _add(srv, payload: Body) -> Body:
scope, actor = _target_scope(payload), _actor_scope(payload)
modality = Modality(payload.get("modality", "text"))
if modality == Modality.VIDEO:
return _submit_video(srv, payload, scope=scope, identity=actor)
# metadata 透传:infer 等调用级开关经 metadata 下推到引擎(engine.write 从
# metadata["infer"]=="true" 判定是否同步走 evolve(EXTRACT) 抽取派生记忆)。
# JSON 标量原样透传(不 str 化):数值/布尔要保持原生类型才能在索引里建
Expand Down Expand Up @@ -613,9 +706,13 @@ def _evolve(srv, payload: Body) -> Body:


def _job(srv, payload: Body) -> Body:
"""查询演进任务状态(Scheduler)。"""
"""查询视频 Ingest 任务或原生 Scheduler 任务状态。"""
job_id = str(_require(payload, "job_id"))
actor = _actor_scope(payload)
info = srv.api.job_status(_require(payload, "job_id"), identity=actor)
if job_id.startswith(INGEST_JOB_PREFIX):
info = srv.ingest_jobs.status(job_id, scope=_target_scope(payload))
return _ingest_job_status(srv, info, identity=actor)
info = srv.api.job_status(job_id, identity=actor)
return {
"ok": True,
"op": "job",
Expand Down
18 changes: 18 additions & 0 deletions bootstrap/core/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,9 @@
Kernel = _api_module.Kernel
build_kernel = _api_module.build_kernel
KernelConfig = import_module("jiuwen_memory.config").Config
IngestJobController = import_module(
"jiuwen_memory.control.ingest_job"
).IngestJobController


class Server:
Expand All @@ -43,6 +46,17 @@ class Server:
def __init__(self, config: Config, kernel: Kernel) -> None:
self.config = config
self.kernel = kernel
memory_config = config.settings.get("memory_api", {})
globals_config = (
memory_config.get("globals", {})
if isinstance(memory_config, dict)
else {}
)
self.ingest_jobs = IngestJobController(
max_workers=int(globals_config.get("ingest_max_workers", 1)),
max_pending_jobs=int(globals_config.get("ingest_max_pending_jobs", 2)),
kv=kernel.kv,
)

@property
def api(self):
Expand Down Expand Up @@ -75,6 +89,10 @@ def dispatch(self, verb: str, payload: Dict[str, Any]) -> Tuple[int, Dict[str, A

return _dispatch(self, verb, payload)

def close(self, *, wait: bool = True) -> None:
"""Release the Control-owned ingest worker pool."""
self.ingest_jobs.close(wait=wait)


def default_spaces() -> Dict[str, Any]:
"""Default scope/namespace registry (none needed for the in-memory build)."""
Expand Down
1 change: 1 addition & 0 deletions bootstrap/http_server/__main__.py
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,7 @@ def serve(self, host: str, port: int) -> None:
sys.stderr.write("\nagent-memory server stopped\n")
finally:
httpd.server_close()
self.close(wait=True)


def main(argv: list[str] | None = None) -> int:
Expand Down
66 changes: 66 additions & 0 deletions examples/config_multimodal.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
memory_api:
globals:
ingest_max_workers: 1
ingest_max_pending_jobs: 2

normalizer:
default:
target: routing
params:
routes:
video:
target: video
params:
whisper_model_dir: /path/to/whisper-model
whisper_batch_size: 1
vllm_base_url: http://127.0.0.1:8000/v1
vllm_api_key: dummy
llm_model: qwen-vl
temp_root: /tmp/agent-memory-video

extractor:
video:
target: video_memory

retriever:
multimodal:
target: multimodal
params:
base_retriever: default
kv_store: default
clip_top_k: 10
event_top_k: 10
rrf_k: 60

evolver:
video:
target: orchestrating
params:
extractor: video
abstractor: default
associator: default
index_builder: default
kv_store: default
graph_store: default
dedup: default
llm: default

pipeline:
default:
target: metadata
params:
route_key: pipeline
fallback: default
routes:
video: video
profiles:
default:
index_builder: default
retriever: multimodal
evolver: default
classifier: default
video:
index_builder: default
retriever: multimodal
evolver: video
classifier: default
2 changes: 2 additions & 0 deletions jiuwen_memory/common/normalizer/normalizer_impl/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,5 +8,7 @@
from jiuwen_memory.common.normalizer.base import NormalizerProducer

import_module(".passthrough_normalizer", __name__)
import_module(".routing_normalizer", __name__)
import_module(".video_normalizer", __name__)

__all__ = ["NormalizerProducer"]
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
"""Route raw payloads to a modality-specific normalizer."""

from __future__ import annotations

from collections.abc import Mapping

from jiuwen_memory.common.base import PluginType
from jiuwen_memory.common.errors import ValidationError
from jiuwen_memory.common.normalizer.base import Normalizer, NormalizerProducer
from jiuwen_memory.common.type_def import Modality, RawPayload


class RoutingNormalizer(Normalizer):
"""Delegate normalization while keeping one Ingestor entry point."""

def __init__(
self,
fallback: Normalizer,
routes: Mapping[Modality, Normalizer],
) -> None:
self._fallback = fallback
self._routes = dict(routes)

def modalities(self) -> list[Modality]:
supported = set(self._fallback.modalities())
supported.update(self._routes)
return sorted(supported, key=lambda item: item.value)

def plugin_type(self) -> PluginType:
return PluginType.NORMALIZER

def health(self) -> None:
seen: set[int] = set()
for normalizer in [self._fallback, *self._routes.values()]:
if id(normalizer) in seen:
continue
seen.add(id(normalizer))
normalizer.health()

def normalize(self, payload: RawPayload) -> str:
normalizer = self._routes.get(payload.modality, self._fallback)
if payload.modality not in normalizer.modalities():
raise ValidationError(
f"no normalizer configured for modality {payload.modality.value!r}"
)
return normalizer.normalize(payload)


def _build_normalizer(config, value: object, *, field: str) -> Normalizer:
if isinstance(value, str):
return NormalizerProducer.build_named(value, config.ctx)
if isinstance(value, Mapping):
target = str(value.get("target", "")).strip()
if not target:
raise ValidationError(f"routing normalizer {field!r} is missing target")
return NormalizerProducer.build(
target,
value.get("params", {}),
config.ctx,
name=str(value.get("name", "")),
)
raise ValidationError(
f"routing normalizer {field!r} must be a named reference or component mapping"
)


@NormalizerProducer.register("routing")
def _build(config):
fallback_raw = config.get("fallback", {"target": "passthrough"})
fallback = _build_normalizer(config, fallback_raw, field="fallback")
routes_raw = config.get("routes", {})
if not isinstance(routes_raw, Mapping):
raise ValidationError("routing normalizer routes must be a mapping")
routes: dict[Modality, Normalizer] = {}
for modality_name, raw in routes_raw.items():
try:
modality = Modality(str(modality_name))
except ValueError as exc:
raise ValidationError(
f"unknown routing normalizer modality {modality_name!r}"
) from exc
routes[modality] = _build_normalizer(
config,
raw,
field=f"routes.{modality.value}",
)
return RoutingNormalizer(fallback, routes)
Loading