Skip to content

Pluggable trajectory processing and runtime profiles

This PR implements processor registration/discovery, immutable execution plans, versioned profiles, and Python/REST/MCP/CLI adapters. Profile processing is opt-in on existing ingestion paths; calls without a profile retain their legacy behavior.

Why profiles are separate from namespaces

Our application selects standard or consistency guidelines per namespace. Other applications select behavior per agent, user, or request. Third-party processing may have nothing to do with guidelines.

A profile selects processor plugins and their configuration. The application decides which profile applies. A namespace remains a persistence destination, not a mandatory configuration ownership model. The core never needs a new field for third-party configuration.

Define processing with real built-in settings

{
  "schema_version": 1,
  "processors": [
    {
      "id": "guidelines",
      "plugin": "evolve.guidelines",
      "config": {
        "guidelines_mode": "consistency",
        "consistency_method": "accurate",
        "segmentation_enabled": false
      }
    }
  ]
}

The built-in adapts the existing standard, fast consistency, and accurate consistency implementations. guidelines_mode accepts standard, consistency, or all; consistency_method accepts fast or accurate. Its schema also includes the model/provider, segmentation, and optional accurate-analysis configuration. Omitted values are resolved when validating/publishing, not read repeatedly during processing. An omitted analysis configuration is captured from the shipped analyzer YAML. An explicit analysis configuration replaces it; it is not a path to a file that can change under a running job. Debug-output location remains deployment-controlled.

GuidelineProcessor.from_config() selects the generation functions and captures their runtime options when each processor instance is constructed. process() only runs those selected steps. A profile update is picked up when the next input batch resolves its plan and constructs fresh processors; in-flight instances and pinned plans keep their original selection.

Each instance ID is unique within the profile. Multiple instances may use the same plugin. The processor list is ordered; an empty list explicitly runs no processors and returns a warning in processing diagnostics. Updates replace the definition, so removed config fields return to plugin defaults. Plugin validation runs before publication. Invalid modes do not silently fall back.

Python: no profile persistence required

from altk_evolve.processing import ProcessingManager

processing = ProcessingManager()  # built-ins + installed entry points; in-memory profiles
plan = processing.validate(definition)
result = processing.process({"messages": messages}, plan=plan)
# result.entities contains proposed entities; this form performs no database writes.

To persist outputs through the existing backend and memory hooks:

from altk_evolve.frontend.client.evolve_client import EvolveClient

client = EvolveClient(processing=processing)
client.ensure_namespace("memories")
result = client.process_trajectory(
    {"messages": messages, "trace_id": "task-123"},
    namespace_id="memories",
    plan=plan,
)

Without an explicit manager, EvolveClient.processing uses the backend's profile repository in the existing configured database. PostgreSQL stores a processing_profiles table alongside the entity tables, using the same connection settings. Filesystem profiles default to entities.sqlite.db inside the configured data directory; its namespace catalog remains in JSON files. EVOLVE_SQLITE_PATH / EVOLVE_SQLITE_URI can explicitly select another file. Milvus profiles use the backend's sqlite_uri, including its existing EVOLVE_SQLITE_PATH override. Use absolute shared paths when clients run from different working directories. Previously published filesystem profiles in a CWD-relative file can be retained by explicitly selecting that file. No separate profile database setting is needed. Repository injection remains available for custom integrations.

Python: saved profiles and application selection

from altk_evolve.processing import ProfileReference

created = client.processing.put("support-review", definition, expected_revision=0)
updated = client.processing.put("support-review", new_definition, expected_revision=created["revision"])
plan = client.processing.resolve("support-review", revision=updated["revision"])
result = client.process_trajectory(trajectory, namespace_id="memories", plan=plan)

Revision 0 means create-only. Subsequent writes require the last observed revision; stale writes raise ProfileConflict. Old revisions remain available. Both SQLite and PostgreSQL use transactions and serialize competing profile writes across repository instances. PostgreSQL uses operation-local connections to the same configured database, so an admin can publish while a trajectory uses its already-captured plan.

An application can choose profiles directly or inject a selector:

client = EvolveClient(
    processing=processing,
    processing_selector=lambda context: ProfileReference(id=context["agent_profile"]),
)
result = client.process_trajectory(
    trajectory,
    namespace_id="memories",
    context={"agent_profile": "support-review"},
)

The selector may return a ProcessingPlan, a profile name, a ProfileReference, or None. None means use the compatibility default built-in plan. An explicit plan or profile bypasses the selector. Unknown references are errors, not default fallback. The context is application-owned; no user/agent hierarchy or authentication is inferred.

Namespace, user, and agent profile bindings belong to the embedding application. For our application's namespace metadata, a binding could be {"processing_profile": {"id": "support-review"}}; this PR does not add a namespace metadata column or a namespace-specific binding API. It provides the selection seam without imposing a storage model on other applications.

Runtime switching and pinned jobs

# Latest: refresh before each trajectory.
for trajectory in trajectories:
    client.process_trajectory(
        trajectory, namespace_id="memories", processing_profile="support-review"
    )

# Pinned: resolve once, including defaults and implementation references.
plan = client.processing.resolve("support-review", revision=1)
for trajectory in trajectories:
    client.process_trajectory(trajectory, namespace_id="memories", plan=plan)

If A starts under revision 1 and an update publishes revision 2, A finishes with revision 1; the next latest-following trajectory uses revision 2. This applies to processor additions/removals as well as settings. Plans capture processor classes in an ordered tuple and store resolved configuration once in the serialized manifest. Each invocation validates an isolated config and calls the class's from_config(config) factory to create a fresh processor instance; plugins receive isolated trajectory copies.

Profiles store normalized defaults and plugin versions. Resolving a saved profile rejects changed processor versions or configuration drift; publish a new revision when processor compatibility changes. The built-in owns version 1, independently of the Evolve package version; ordinary package releases do not invalidate profiles. A retained in-process plan keeps its processor class references. Hot replacement of installed Python code is unsupported.

Built-ins, discovered packages, and local plugins

A processor class has id, api_version=1, version, a Pydantic config_model, a from_config(config) classmethod that constructs its configured instance, and a process(trajectory, *, context) instance method returning ProcessorResult. Config schemas are plugin-owned. The result contains entities, diagnostics, and an optional request for persistence-time conflict resolution. No subclass is required. Registration and inventory inspect class metadata and never construct instances. Construction belongs to the plugin, so its constructor can require configuration or plugin-specific dependencies:

class MyProcessor:
    id = "example.custom"
    api_version = 1
    version = "1.0"
    config_model = MyConfig

    @classmethod
    def from_config(cls, config):
        return cls(MyConfig.model_validate(config))

    def __init__(self, config):
        self.config = config

    def process(self, trajectory, *, context):
        return ProcessorResult(entities=[])

MyConfig is the application's Pydantic model. The runner supplies a fresh validated config to from_config for every trajectory, including repeated runs of a pinned plan. There is no separate BoundProcessor record or external factory callable.

Built-ins are registered automatically. Installed packages advertise entry points:

[project.entry-points."altk_evolve.processors"]
"example.word_count" = "word_count:WordCountProcessor"

Discovery enumerates registrations without running processors. Schemas/implementations are loaded on demand. This uses the standard Python entry-point mechanism. Duplicate IDs (including built-in shadowing), incompatible APIs, and missing selected plugins fail explicitly. Inventory reports unavailable plugins. Restart after installing new packages. Python plugins execute trusted code in the host; profiles cannot install packages or supply arbitrary import paths.

Local plugins use the same registry:

from altk_evolve.processing import ProcessorRegistry, ProcessingManager

registry = ProcessorRegistry.discover()
registry.register(MyProcessor)
processing = ProcessingManager(registry=registry)

Use discover(installed=False) to exclude installed extensions, or ProcessorRegistry() for an empty registry. Registration makes a processor available; only a plan/profile activates it.

A complete no-LLM package example is in examples/processing_plugin. Install it into the same environment as the CLI, then run:

uv pip install --no-deps -e examples/processing_plugin
uv run evolve processors list
uv run evolve processing-profiles apply word-count --file examples/processing_plugin/profile.json --expected-revision 0
uv run evolve namespaces create memories
uv run evolve processing run --file examples/processing_plugin/trajectory.json --namespace memories --processing-profile word-count

REST and MCP

REST is mounted under /api on the existing FastAPI service:

Action REST MCP tool
IDs, versions, JSON config schemas GET /api/processors list_processors
Read profile/latest or pinned revision GET /api/processing-profiles/{id}?revision=1 get_processing_profile(profile_id, revision=None)
Validate and replace profile PUT /api/processing-profiles/{id} set_processing_profile(profile_id, definition, expected_revision)
Run processors and persist derived entities POST /api/trajectories process_trajectory(trajectory, namespace_id, processing_profile, revision=None)

PUT takes the complete definition above. Create with If-None-Match: *; update with If-Match: "1". Responses include ETag and revision. Missing write preconditions return 428, stale writes 409, missing profiles 404, and invalid config 422. MCP uses expected_revision=0 for create-only.

POST body:

{
  "namespace_id": "memories",
  "processing_profile": {"id": "support-review", "revision": 1},
  "trajectory": {
    "trace_id": "task-123",
    "messages": [{"role": "user", "content": "Review this task"}]
  }
}

Omit revision to select latest. The destination namespace must already exist. These execution APIs persist derived entities, not the raw trajectory. Existing MCP save_trajectory also accepts processing_profile and optional profile_revision; it preserves its raw-trajectory persistence and return shape. The service's existing authentication model is unchanged. Application adapters must supply their own authorization; identifiers alone are not authenticated principals. The stock REST/MCP service has no authenticated-principal or profile-editor policy. Run it within a trusted deployment boundary, or protect all routes with application authentication and resource authorization before exposing it to untrusted callers.

CLI and Phoenix sync

CLI profile operations connect to the configured database, including PostgreSQL. They do not contact a running REST/MCP service. Clients using the same database see the same profile revisions; use REST/MCP when the host service owns the database access.

uv run evolve processing-profiles get support-review --revision 1
uv run evolve sync phoenix --processing-profile support-review
uv run evolve sync phoenix --processing-profile support-review --profile-revision 1

Without a revision, Phoenix sync resolves before each completed LLM span; a pinned revision is resolved once when the syncer is constructed. A profile cannot be combined with legacy --guidelines-mode or --consistency-method flags. Existing invocations without a profile remain compatible. Profiles are activated explicitly, so installing an extension does not change existing sync behavior.

Persistence, hooks, and failure behavior

Processors execute sequentially against the same original input, without consuming each other's results. All must succeed before derived writes begin. Failures stop the profile operation rather than silently discarding a processor's output. This differs intentionally from the legacy MCP save_trajectory path, which catches generation errors and skips failed standard/consistency generation. Profile processing propagates those errors, including transient LLM failures, so the caller can retry.

Persistence groups outputs by entity type and uses the normal backend path. Existing memory hooks still run; built-in generation preserves LLM-egress hooks. Third parties can use context.complete(...) for mediated LLM calls. Arbitrary direct network calls from trusted third-party Python cannot be intercepted by this contract.

The runner captures conflict-resolution model/provider settings with the plan and passes them through the client/backend/LLM call chain. Generation receives explicit settings; it does not mutate globals to switch modes. Per-call LiteLLM global JSON validation toggles were replaced with per-request validation in these generation paths; responses are also locally validated with Pydantic. Deployment connections, credentials, and global hook setup are not hot-swapped by profiles. Constructing clients with different hook configurations still has the existing process-global hook lifecycle limitation.

Every produced entity carries the full effective manifest, digest, operation ID, processor instance ID, and optional profile ID/revision. This duplicates small manifests rather than requiring a separate artifact store. Provenance is stamped after conflict resolution so model-returned metadata cannot replace it. Unchanged entities retain prior provenance. Results expose proposed entities and actual persistence updates separately; an update may consolidate into an existing entity.

Incremental input and checkpoints

A conversation can remain open indefinitely. Supply Trajectory.batch to identify a bounded contribution; messages contains new material and context_messages contains supporting history. For example, the same chat can submit events 1–120, then 121–140:

result = client.process_trajectory(
    {
        "batch": {
            "source": "my-app",
            "conversation_id": "chat-42",
            "batch_id": "events-121-140",
            "revision": "1",
            "scope": "agent:researcher",
        },
        "messages": new_messages,
        "context_messages": earlier_messages,
    },
    namespace_id="memories",
    processing_profile="support-review",
)

REST's trajectory object, MCP process_trajectory, and the CLI input JSON accept the same fields. The source adapter owns stable, non-overlapping batch IDs and changes revision when that input is corrected. Message count is not a source identity. Scope is application-owned: choose it to separate agent/user consumers sharing a namespace, or keep the default for namespace-wide processing. Use a new scope for an explicit replay. Changing a profile revision does not replay completed contributions. Without batch, calls remain untracked and repeated calls may generate new outputs.

Each processor instance has independent progress, keyed by namespace, scope, source, conversation, processor ID, batch ID, and source revision. Processing captures the configuration once per invocation. After generation, hooks, reconciliation and embedding computation, a short commit stores that processor's outputs and checkpoint together. Completed processors are skipped on redelivery. If persistence fails for a later processor, earlier commits stay complete; a subsequent invocation processes only the unfinished instances. Its selected profile is captured for that invocation. completed_processors and skipped_processors distinguish committed work from proposals in entities; updates reports actual storage mutations.

Unrelated namespace changes do not invalidate processing. Reads for semantic reconciliation can be slightly stale. Only replacement/deletion targets are compared before destructive writes on atomic backends. Advisory last_accessed changes are rebased onto updates; content and policy-relevant metadata changes still raise ConcurrentEntityUpdate, leaving that contribution uncommitted and eligible for redelivery. There is no namespace-wide validation or automatic model retry. Concurrent deliveries can both run a processor, but only one can commit the same checkpoint.

Profile-enabled Phoenix sync processes each completed innermost LLM span independently, using its span ID and a hash of its processing payload as batch identity and revision. The completion is new material and the prompt is supporting context. This handles late spans, corrected content, and additional calls in the same trace without declaring the trace permanently done. Running spans wait until end_time is present. Each poll still observes only the configured fetched span window (limit); configure it to cover the arrival rate. The legacy no-profile sync path is unchanged. Profile sync leaves raw transcripts in Phoenix rather than creating a trajectory entity as a completion marker.

Filesystem stores checkpoints in the namespace JSON, publishing output and progress with one atomic rename under its existing process-shared writer lock. PostgreSQL stores processing_checkpoints(namespace_id, key, value) alongside entity tables in the same database, committing both on a dedicated connection. No temporary table or whole-namespace snapshot is taken during generation. The backend transaction is a storage-only boundary; call processing outside it.

Hooks receive HookBackend for reads and metadata-patch proposals. During preparation, patches are collected and applied with the output commit; no hook is invoked under the commit lock. prepare_updates() returns an opaque receipt bound to its backend, namespace, and checked changes; constructing, modifying, or retargeting it is rejected by commit_prepared(). Multiple prepared mutations of the same target are rejected. Callbacks must finish before returning. External side effects are outside the storage guarantee and must tolerate duplicate execution.

An atomic backend declares supports_atomic_writes and implements transaction, get_processing_checkpoint, and _save_processing_checkpoint. Checkpoint reads are internal storage operations, independent of entity hooks. Milvus and other backends without this capability reject identified batches before processor execution; untracked processing remains available with its existing non-atomic storage behavior.

Validation and remaining scope

Unit tests cover revision conflicts, restart persistence, plugin validation/discovery, mode switching during processing, pinned plans, isolated configs, conflict provenance, application selectors, shared transport APIs, and existing ingestion compatibility. A subprocess E2E test discovers a distribution entry point, updates profiles through separate CLI invocations, and verifies pinned/latest results in filesystem storage. Existing guideline, consistency, backend, hook, and CLI tests remain regression gates.

Not implemented: namespace binding storage, remote CLI management, execution DAGs, parallel processors, plugin hot code reload, and automatic retries/cancellation. SQLite and PostgreSQL profile publication use atomic revision checks. A caller may inject another profile repository; none of the processing interfaces require namespace-specific SQL or a fixed user/agent model.

Accurate consistency configuration validation

Partial analysis_config dictionaries merge over the bundled analyzer defaults, so {} preserves its agent metrics and sampling behavior. Profile publication checks finite numbers, the uncertainty threshold range (0, 1], aggregation, boolean skip behavior, and agent descriptors. Profile sampling is limited to 1–100 samples and 1–1000 steps to bound each invocation; these are profile API limits. Analyzer-specific metric configuration remains extensible. Configurations must also survive JSON serialization and revalidation identically before a revision is written.

Conflict-resolution updates keep the latest processing stamp and append the prior stamp to processing_history, sourced from stored entities rather than LLM output. History retains complete manifests so unpublished/ad-hoc plans remain traceable.

Accurate consistency preserves context_messages as the original structured inference prefix, including system instructions and tool exchanges. Only assistant turns in messages are numbered and resampled. Guideline generation receives historical context separately from those scored steps. All three guideline paths bound the history rendered for extraction to 20,000 characters (up to 50 recent messages, 2,000 characters each), and retain the original task separately. This limit never alters the replay prefix.

Phoenix records observed ancestor LLM span IDs as batch aliases, so discovering an inner instrumentation span later does not repeat an already committed contribution with the same conversation content. Each alias carries its own representation revision, so provider/model labels and tool-schema detail need not be byte-identical. Changed conversation content is not aliased. Alias links commit atomically and survive subsequent polls and restarts. Discovery requires the span ancestry to be present in the fetched window; configure the fetch limit accordingly. Independent calls with identical content remain independent batches.

Refreshing application-owned defaults

Applications can use ProcessingManager.ensure(name, definition) or the MCP tool ensure_processing_profile(profile_id, definition) for a profile they own. Defaults resolve on the Evolve service. The operation keeps the profile ID, creates an immutable revision only when needed, and retries conditional-write conflicts. Pass the returned revision when processing a batch.

Evolve retains the previously resolved defaults alongside the profile. Fields that still match those defaults follow updated configuration; changed fields, added processors, and operator removals are preserved. Explicit values in the application's supplied definition remain explicit. Regular profile updates keep the default baseline so later refreshes preserve those edits. Historical profile revisions continue to resolve to their original models and providers.

When adopting an existing profile, its first revision serves as the baseline. Use this only for an application-owned profile whose first revision represents that application's defaults, not an unrelated administrator-owned profile.