How to Customize the Memorize Pipeline in memU: Insert, Remove, or Replace Workflow Steps
You customize the memorize pipeline in NevaMind-AI/memU by calling MemoryService helper methods—insert_step_before, insert_step_after, replace_step, or remove_step—which delegate to PipelineManager to create immutable revisions while preserving running workflows.
The memU library provides a flexible, revision-aware architecture for building memory extraction workflows. Whether you need to add custom preprocessing, swap out the default item extractor, or remove unnecessary deduplication stages, the pipeline modification API in src/memu/app/service.py exposes high-level methods that safely mutate the memorize workflow without breaking in-flight processes.
Understanding the Memorize Pipeline Architecture
Before modifying steps, you need to understand three core components that work together in NevaMind-AI/memU:
WorkflowStep– Defined insrc/memu/workflow/step.py, this dataclass bundles astep_id,role,handlercallable, and metadata fields (requires,produces,capabilities,config) that declare dependencies and outputs.PipelineManager– Located insrc/memu/workflow/pipeline.py, this class maintains a dictionary of pipeline names mapped to lists ofPipelineRevisionobjects. Every mutation creates a new revision, ensuring that running workflows use the snapshot they started with while new calls pick up the latest version.MemoryService– The public API insrc/memu/app/service.pythat constructs the default memorize workflow and exposes convenience methods (insert_step_before,replace_step, etc.) which forward to thePipelineManager.
How to Insert Workflow Steps into the Memorize Pipeline
Adding new processing stages is the most common customization. The MemoryService provides two insertion points relative to an existing step ID.
Inserting a Step Before an Existing Step
Use insert_step_before when you need to preprocess data before a specific stage runs. For example, inserting a custom OCR step before preprocess_multimodal:
from memu.app.service import MemoryService
from memu.workflow.step import WorkflowStep, WorkflowState, WorkflowContext
async def custom_ocr_handler(state: WorkflowState, ctx: WorkflowContext) -> WorkflowState:
# Custom OCR logic here
state["ocr_text"] = "Extracted text"
return state
ocr_step = WorkflowStep(
step_id="custom_ocr",
role="ocr",
handler=custom_ocr_handler,
requires={"local_path"},
produces={"ocr_text"},
capabilities={"io"},
)
service = MemoryService()
revision = service.insert_step_before(
target_step_id="preprocess_multimodal",
new_step=ocr_step,
pipeline="memorize",
)
print(f"Inserted step at revision {revision}")
Inserting a Step After an Existing Step
Use insert_step_after to append processing logic following a specific step. This is useful for adding post-processing tags after categorize_items:
revision = service.insert_step_after(
target_step_id="categorize_items",
new_step=tagging_step,
pipeline="memorize",
)
How to Replace or Remove Steps in the Memorize Pipeline
When the default workflow contains steps that do not fit your domain, you can swap them out or delete them entirely.
Replacing an Existing Step
The replace_step method substitutes a step while preserving the pipeline order. This is ideal when you need a specialized extractor instead of the default extract_items step:
async def specialized_extractor(state: WorkflowState, ctx: WorkflowContext) -> WorkflowState:
# Domain-specific extraction logic
state["resource_plans"] = [{"entries": []}]
return state
new_extract_step = WorkflowStep(
step_id="extract_items", # Keep same ID to maintain downstream compatibility
role="extract",
handler=specialized_extractor,
requires={"preprocessed_resources", "memory_types", "categories_prompt_str", "modality", "resource_url"},
produces={"resource_plans"},
capabilities={"llm"},
)
new_rev = service.replace_step(
target_step_id="extract_items",
new_step=new_extract_step,
pipeline="memorize",
)
Removing an Unnecessary Step
Use remove_step to delete stages like dedupe_merge if your use case does not require deduplication:
rev = service.remove_step(target_step_id="dedupe_merge", pipeline="memorize")
print(f"Removed step – pipeline now at revision {rev}")
Configuring Step Parameters Without Code Changes
Not every customization requires writing a new handler. The configure_pipeline method updates the configuration dictionary of an existing step, allowing you to adjust LLM profiles, prompts, or timeouts without replacing the step object:
rev = service.configure_pipeline(
step_id="preprocess_multimodal",
configs={"chat_llm_profile": "fast-gpt-4", "timeout_seconds": 30},
pipeline="memorize",
)
This modifies the config field of the WorkflowStep in the next pipeline revision while keeping the same handler and dependencies.
Complete Working Example: Custom OCR Integration
Here is a full example that combines insertion and configuration to add custom OCR processing before the multimodal preprocessing stage:
import asyncio
from memu.app.service import MemoryService
from memu.workflow.step import WorkflowStep, WorkflowState, WorkflowContext
async def main():
# Define custom OCR handler
async def ocr_handler(state: WorkflowState, ctx: WorkflowContext) -> WorkflowState:
# Simulate OCR processing
state["ocr_text"] = f"OCR result for {state.get('local_path', 'unknown')}"
return state
# Create the step
ocr_step = WorkflowStep(
step_id="custom_ocr_preprocessor",
role="ocr",
handler=ocr_handler,
requires={"local_path"},
produces={"ocr_text"},
capabilities={"io"},
description="Custom OCR preprocessing",
)
# Initialize service and modify pipeline
service = MemoryService()
# Insert before preprocess_multimodal
revision = service.insert_step_before(
target_step_id="preprocess_multimodal",
new_step=ocr_step,
pipeline="memorize",
)
print(f"Pipeline customized at revision {revision}")
# Run the modified workflow
result = await service.memorize(
resource_url="https://example.com/document.pdf",
modality="document",
user={"user_id": "alice"},
)
print(result)
if __name__ == "__main__":
asyncio.run(main())
Summary
- MemU’s memorize pipeline is a mutable sequence of
WorkflowStepobjects managed byPipelineManagerand exposed throughMemoryService. - Insert steps using
insert_step_beforeorinsert_step_afterto add preprocessing or post-processing stages without disrupting existing logic. - Replace steps using
replace_stepwhen you need to swap out handlers likeextract_itemswhile preserving pipeline order and downstream compatibility. - Remove steps using
remove_stepto eliminate unnecessary stages such asdedupe_merge. - Configure steps using
configure_pipelineto adjust runtime parameters like LLM profiles without writing new handler code. - All mutations are revision-aware, creating immutable snapshots in
PipelineManagerso running workflows remain stable while new executions pick up the latest configuration.
Frequently Asked Questions
How do I ensure my custom step receives the correct input data?
Define the requires set in your WorkflowStep to match keys produced by previous steps. The PipelineManager validates dependencies between revisions, ensuring that requires fields are satisfied by earlier steps' produces sets. For example, if your step needs OCR text, set requires={"ocr_text"} and ensure a preceding step sets produces={"ocr_text"}.
Can I modify the memorize pipeline while workflows are running?
Yes. The PipelineManager implements immutable revisions—each modification creates a new PipelineRevision object. Active workflow executions continue using the revision they started with, while subsequent calls to service.memorize() automatically use the latest revision. This design prevents runtime errors and ensures thread-safe updates.
What is the difference between replacing a step and configuring a step?
Replacing a step (replace_step) substitutes the entire WorkflowStep object, including its handler function, role, and dependency declarations. Use this when you need completely different logic. Configuring a step (configure_pipeline) updates only the config dictionary of an existing step, allowing you to change runtime parameters like LLM profiles, temperature settings, or prompt templates without altering the handler code.
How do I revert to a previous pipeline revision?
Currently, the MemoryService API exposes forward-only mutations (insert, replace, remove, configure). To revert, you must re-apply the desired state by calling the appropriate modification methods to reconstruct the previous configuration. For programmatic rollback support, you can access the PipelineManager directly via service._pipeline_manager and inspect self._pipelines["memorize"] to retrieve historical revision objects, though this requires internal API usage.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →