How memU's Continuous Learning Pipeline Processes Inputs in Real-Time
MemU's continuous learning pipeline processes inputs in real-time through an asynchronous, contract-based workflow system that converts every incoming datum into a persisted memory item via the PipelineManager, LocalWorkflowRunner, and MemoryService components without interrupting the running application.
The NevaMind-AI/memU repository implements an as-you-type continuous learning architecture that treats every incoming piece of information as a discrete workflow. This continuous learning pipeline enables real-time memory formation by processing user messages, documents, and tool outputs through an immutable, versioned pipeline system that never blocks the main application thread.
Core Architecture of the Continuous Learning Pipeline
The pipeline rests on three foundational components that handle revision management, execution, and public API exposure.
PipelineManager and Immutable Revisions
The PipelineManager class in src/memu/workflow/pipeline.py stores immutable pipeline revisions, ensuring that mutations create fresh PipelineRevision instances rather than altering running workflows. When execution begins, PipelineManager.build("memorize") returns a shallow copy of the current step list (lines 47-50), guaranteeing that each workflow instance receives an isolated snapshot while the underlying revision history remains immutable.
LocalWorkflowRunner and Step Execution
The LocalWorkflowRunner in src/memu/workflow/runner.py (lines 31-40) serves as the default execution engine. It drives pipeline steps sequentially while respecting the asynchronous nature of modern I/O operations. The runner transparently handles both synchronous and asynchronous step implementations, automatically awaiting coroutines when necessary. This async-first design allows the continuous learning pipeline to perform LLM embedding calls and vector database writes without blocking the application thread.
MemoryService as the Public Facade
The MemoryService class in src/memu/app/service.py exposes the continuous learning capability through the HTTP endpoint POST /api/v3/memory/memorize (documented in readme/README_en.md). This service registers all pipelines via _register_pipelines() (lines 15-19) and launches the "memorize" workflow through _run_workflow() (lines 50-55). The service manages WorkflowState initialization, lazy-loads LLM clients via _get_llm_client (lines 87-90), and registers interceptors for cross-cutting concerns like tracing and logging (lines 58-70).
Real-Time Input Processing Flow
The continuous learning pipeline processes inputs through seven distinct phases, each optimized for minimal latency and maximum reliability.
Step 1: HTTP Request Ingestion
Real-time processing initiates when a client POSTs a payload to /api/v3/memory/memorize. The MemoryService routes this request to _run_workflow with the workflow name "memorize", triggering the asynchronous execution chain.
Step 2: WorkflowState Initialization
Inside _run_workflow() (lines 50-55 of src/memu/app/service.py), the JSON body transforms into a WorkflowState dictionary (dict[str, Any]) containing the raw input, a unique trace ID, and user-provided configuration overrides. This state object serves as the shared data context passed between all pipeline steps.
Step 3: Pipeline Resolution and Copying
The PipelineManager.build("memorize") method (lines 47-50 of src/memu/workflow/pipeline.py) resolves the current pipeline definition and returns a shallow copy of the step list. This copy-on-read mechanism ensures that mutations to the pipeline configuration generate new immutable revisions without affecting the currently executing workflow instance.
Step 4: Asynchronous Step Execution
The LocalWorkflowRunner.run() method (lines 31-40 of src/memu/workflow/runner.py) forwards the copied steps to the run_steps engine. Each step executes asynchronously, allowing the continuous learning pipeline to handle I/O-bound operations—such as LLM embedding generation and vector database writes—concurrently with application processing.
Step 5: Contract Validation and State Updates
Before each step executes, the run_steps function (defined in src/memu/workflow/step.py) validates that all keys listed in step.requires exist in the current WorkflowState (lines 69-73). Missing keys raise a KeyError, enforcing strict data flow contracts. After successful execution, the step's produces keys merge into the state (line 90), making the output available to downstream steps.
Step 6: Real-Time Memory Storage
The "memorize" workflow typically contains three specialized steps that execute in immediate succession:
- Extract text:
requires: {}/produces: {"raw_text"} - Embed: Calls the selected LLM embedding client using
requires: {"raw_text"}/produces: {"embedding"} - Store: Writes to the vector database and metadata store via
requires: {"embedding"}/produces: {"memory_id"}
Because each step runs immediately after the previous one completes, the system ingests and indexes the input in real-time, making the new memory searchable instantly. The concrete step definitions reside in _build_memorize_workflow() within MemoryService (around lines 15-19 of src/memu/app/service.py).
Step 7: Response Confirmation
Once all steps succeed, _run_workflow returns the final WorkflowState containing the memory_id, timestamps, and metadata. The MemoryService serializes this state into JSON and sends it back to the caller, confirming that the continuous-learning item has been persisted.
Key Architectural Concepts
Several design patterns enable memU's real-time continuous learning capabilities.
Pipeline Revisioning
Every mutation to a pipeline configuration—such as config_step or insert_after—creates a new immutable PipelineRevision rather than modifying existing definitions. This enables A/B testing of pipeline changes without disrupting already-running continuous learning instances. The PipelineRevision class is defined in src/memu/workflow/pipeline.py (lines 12-18).
Declarative Step Contracts
Each WorkflowStep explicitly declares requires and produces sets. The run_steps engine validates these contracts at runtime, ensuring that data dependencies are satisfied before execution and that outputs are properly propagated through the continuous learning pipeline. This is implemented in src/memu/workflow/step.py (lines 16-26).
Interceptor Pattern
Hooks registered via MemoryService.intercept_before_workflow_step (lines 58-70 of src/memu/app/service.py) run automatically before, after, or on error of any step. These interceptors handle cross-cutting concerns like distributed tracing, logging, and metrics collection without cluttering the core business logic.
Async-First Execution
The LocalWorkflowRunner transparently handles both synchronous and asynchronous step implementations. When a step returns a coroutine, the runner awaits it automatically, allowing the continuous learning pipeline to perform I/O operations—such as LLM embedding calls and vector database writes—without blocking the application thread. This logic resides in WorkflowStep.run() (lines 40-46) and run_steps (lines 90-96).
LLM Client Lazy Loading
To minimize resource consumption, the LLM client initializes only when the first step requiring embeddings executes. The _get_llm_client method (lines 87-90 of src/memu/app/service.py) implements this lazy-loading pattern, keeping the continuous learning pipeline lightweight until actual embedding computation is required.
Practical Implementation Example
The following example demonstrates how to invoke memU's continuous learning pipeline programmatically:
from memu.app.service import MemoryService
# Initialise the service (uses default config files)
mem_service = MemoryService()
# Real-time continuous learning: feed a document
payload = {
"content": "MemU can continuously learn from new data.",
"metadata": {"source": "demo"},
# optional per-step overrides
"config": {"embed_llm_profile": "default"},
}
# The call runs the "memorize" pipeline and returns the new memory ID
result = await mem_service._run_workflow(
workflow_name="memorize",
initial_state={"input": payload},
)
print("New memory stored with ID:", result["memory_id"])
This snippet mirrors what the HTTP endpoint does under the hood—it builds a fresh state, copies the pipeline, runs the steps asynchronously, and finally returns the enriched state containing the new memory identifier.
Core Files and Their Roles
Understanding the continuous learning pipeline requires familiarity with these key source files:
-
src/memu/workflow/pipeline.py– DefinesPipelineManager,PipelineRevision, and the immutable revision handling logic. This file manages pipeline mutations and the shallow-copy mechanism that isolates running workflows from configuration changes. -
src/memu/workflow/step.py– Contains theWorkflowStepdataclass and therun_stepsexecution engine. This module enforces therequires/producescontracts and handles the async/await logic for step execution. -
src/memu/workflow/runner.py– ImplementsLocalWorkflowRunner, the default executor that drives the step-by-step execution. It coordinates withrun_stepsand manages the runner resolution logic for pluggable backends. -
src/memu/app/service.py– The public façade exposing the continuous learning API.MemoryServiceregisters pipelines (including "memorize"), manages LLM client lazy loading via_get_llm_client, and launches workflows through_run_workflow. It also handles interceptor registration for cross-cutting concerns. -
readme/README_en.md– Documents the public HTTP API (POST /api/v3/memory/memorize) and provides usage examples for the continuous learning endpoint.
Summary
MemU's continuous learning pipeline delivers real-time memory formation through these key mechanisms:
-
Immutable Pipeline Revisions: The
PipelineManagercreates shallow copies of pipeline definitions for each execution, ensuring that configuration changes don't affect running workflows while enabling A/B testing of pipeline variations. -
Declarative Step Contracts: Each
WorkflowStepexplicitly declaresrequiresandproducessets, with therun_stepsengine enforcing these contracts at runtime to guarantee data flow correctness. -
Async-First Execution: The
LocalWorkflowRunnerhandles both sync and async steps transparently, allowing I/O-bound operations like LLM embedding calls and vector database writes to execute without blocking the application. -
Lazy Resource Initialization: The
MemoryServicedefers LLM client creation until the first embedding step executes, keeping the pipeline lightweight until actual computation is required. -
Real-Time Storage: The "memorize" workflow processes inputs through extract, embed, and store steps in immediate succession, making new memories searchable instantly via the
POST /api/v3/memory/memorizeendpoint.
Frequently Asked Questions
What makes memU's continuous learning pipeline 'real-time'?
The pipeline achieves real-time processing by treating each input as an independent, asynchronous workflow executed by the LocalWorkflowRunner. Rather than batching inputs or pausing for synchronous I/O, the "memorize" workflow runs extract, embed, and store steps immediately upon receiving a POST /api/v3/memory/memorize request. Because steps execute asynchronously and the MemoryService returns the memory_id only after successful persistence, the system ingests and indexes data instantly while maintaining non-blocking operation.
How does memU prevent pipeline configuration changes from breaking running workflows?
MemU implements immutable pipeline revisions through the PipelineManager class in src/memu/workflow/pipeline.py. When a workflow starts, PipelineManager.build() creates a shallow copy of the current step list (lines 47-50), ensuring that the executing workflow references an isolated snapshot. Any mutations to the pipeline configuration—such as config_step or insert_after—generate new PipelineRevision instances rather than modifying existing ones. This copy-on-read mechanism allows developers to A/B test pipeline changes without disrupting already-running continuous learning instances.
What happens if a step in the continuous learning pipeline fails?
The run_steps engine in src/memu/workflow/step.py enforces strict contract validation before executing each step, checking that all requires keys exist in the current WorkflowState (lines 69-73). Missing keys raise a KeyError, enforcing strict data flow contracts. Additionally, the MemoryService registers interceptors via intercept_before_workflow_step (lines 58-70 of src/memu/app/service.py) that can handle errors, log failures, or trigger recovery logic. Because the pipeline operates on isolated WorkflowState instances, failures in one execution do not corrupt the global state or affect other running workflows.
Can the continuous learning pipeline handle different types of inputs simultaneously?
Yes, the pipeline's declarative step contracts and state-driven architecture make it agnostic to input types. The WorkflowState (defined as dict[str, Any]) can carry arbitrary payloads, and individual steps declare their requires and produces sets to handle specific data types. For example, the "memorize" workflow can process user messages, documents, or tool outputs interchangeably because the extract step normalizes inputs into raw_text, which subsequent embedding and storage steps consume. The MemoryService accepts per-step configuration overrides via the config field in the request payload, allowing different input types to trigger different LLM profiles or storage backends within the same running pipeline instance.
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 →