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 in src/memu/workflow/step.py, this dataclass bundles a step_id, role, handler callable, and metadata fields (requires, produces, capabilities, config) that declare dependencies and outputs.
  • PipelineManager – Located in src/memu/workflow/pipeline.py, this class maintains a dictionary of pipeline names mapped to lists of PipelineRevision objects. 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 in src/memu/app/service.py that constructs the default memorize workflow and exposes convenience methods (insert_step_before, replace_step, etc.) which forward to the PipelineManager.

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 WorkflowStep objects managed by PipelineManager and exposed through MemoryService.
  • Insert steps using insert_step_before or insert_step_after to add preprocessing or post-processing stages without disrupting existing logic.
  • Replace steps using replace_step when you need to swap out handlers like extract_items while preserving pipeline order and downstream compatibility.
  • Remove steps using remove_step to eliminate unnecessary stages such as dedupe_merge.
  • Configure steps using configure_pipeline to adjust runtime parameters like LLM profiles without writing new handler code.
  • All mutations are revision-aware, creating immutable snapshots in PipelineManager so 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →