> ## Documentation Index
> Fetch the complete documentation index at: https://docs.evermind.ai/llms.txt
> Use this file to discover all available pages before exploring further.

# Batch Processing

> Import existing conversation history into EverOS at scale

When you have existing chat history to import (Slack exports, support logs, past
sessions), batch processing ingests it efficiently. This guide covers the data
format, a sync importer that parallelizes across conversations, retries, and
resumable checkpoints.

## Prerequisites

```bash theme={null}
pip install everos-cloud
export EVEROS_API_KEY="your_api_key"
```

## Data format

Everything is a **session** of messages. Each message names its speaker with
`sender_id`. A one-on-one conversation uses the user's id plus `"assistant"`; a
group conversation uses several ids. Same format either way.

```json theme={null}
{
  "session_id": "slack_engineering_2026_07",
  "messages": [
    {"sender_id": "user_alice", "role": "user", "timestamp": 1705312800000,
     "content": "Let's discuss the new architecture proposal."},
    {"sender_id": "user_bob", "role": "user", "timestamp": 1705312890000,
     "content": "I reviewed it — I think we should consider microservices."}
  ]
}
```

| Field | Required | Description |
| - | - | - |
| `sender_id` | Yes | Who sent the message (the participant it's attributed to) |
| `role` | Yes | `"user"` or `"assistant"` |
| `timestamp` | Recommended | Unix **milliseconds**. Set it when importing history so ordering and time-based recall are correct |
| `content` | Yes | Message text (or a content array for multimodal) |

<Warning>
  `timestamp` must be unix **milliseconds** (`>= 1_000_000_000_000`). Seconds are
  the most common cause of a `422` on import. Multiply by 1000.
</Warning>

## Convert your data

Map your source export into the session format. Sort chronologically, because
EverOS uses timestamps for boundary detection.

```python theme={null}
from datetime import datetime

def convert_slack_export(slack_messages: list, channel: str) -> dict:
    messages = [
        {
            "sender_id": m.get("user", "unknown"),
            "role": "user",
            "timestamp": int(float(m.get("ts", 0)) * 1000),  # Slack ts is seconds
            "content": m.get("text", ""),
        }
        for m in slack_messages
        if m.get("text")
    ]
    messages.sort(key=lambda m: m["timestamp"])
    return {"session_id": f"slack_{channel}", "messages": messages}
```

## Batch importer

The importer sends each conversation in chunks (the `add` endpoint takes up to
500 messages per call), retries transient failures, and flushes at the end so
extraction starts promptly. Conversations are imported in parallel with a thread
pool; chunks *within* a conversation go in order to preserve the timeline.

```python theme={null}
import json, time, logging
from pathlib import Path
from concurrent.futures import ThreadPoolExecutor, as_completed

from everos_cloud import EverOS, EverOSAPIError

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

client = EverOS()  # one shared, thread-safe client

CHUNK_SIZE = 200          # <= 500 (server limit)
TRANSIENT = {429, 500, 502, 503, 504}


def _add_with_retry(session_id: str, chunk: list, attempts: int = 4):
    for attempt in range(attempts):
        try:
            return client.add(session_id=session_id, messages=chunk)
        except EverOSAPIError as e:
            if e.status not in TRANSIENT or attempt == attempts - 1:
                raise
            time.sleep(0.5 * (2 ** attempt))


def import_conversation(data: dict) -> dict:
    session_id = data["session_id"]
    messages = data["messages"]
    sent = failed = 0

    for i in range(0, len(messages), CHUNK_SIZE):
        chunk = messages[i:i + CHUNK_SIZE]
        try:
            _add_with_retry(session_id, chunk)   # chunks in order
            sent += len(chunk)
        except EverOSAPIError as e:
            logger.error("chunk failed for %s: %s", session_id, e.body)
            failed += len(chunk)

    client.flush(session_id)  # kick off extraction now (also runs on its own schedule)
    logger.info("%s: %d sent, %d failed", session_id, sent, failed)
    return {"session_id": session_id, "sent": sent, "failed": failed}


def import_directory(directory: str, max_workers: int = 8) -> list:
    files = sorted(Path(directory).glob("*.json"))
    logger.info("importing %d conversations", len(files))

    def run(path):
        return import_conversation(json.loads(path.read_text()))

    results = []
    with ThreadPoolExecutor(max_workers=max_workers) as pool:
        futures = {pool.submit(run, p): p for p in files}
        for fut in as_completed(futures):
            try:
                results.append(fut.result())
            except Exception as e:
                logger.error("failed %s: %s", futures[fut], e)
                results.append({"file": str(futures[fut]), "error": str(e)})
    return results


if __name__ == "__main__":
    print(import_directory("./exports/"))
```

<Note>
  Writes are asynchronous server-side. Add returns as soon as messages are
  accepted, and extraction runs in the background. The `flush` at the end
  accelerates extraction for that session; it isn't required for the data to be
  processed.
</Note>

## Resumable imports

For large jobs, checkpoint completed sessions so a re-run skips them.

```python theme={null}
import json
from pathlib import Path

class Checkpoint:
    def __init__(self, path: str = "import_checkpoint.json"):
        self.path = Path(path)
        self.done = set(json.loads(self.path.read_text())) if self.path.exists() else set()

    def is_done(self, session_id: str) -> bool:
        return session_id in self.done

    def mark(self, session_id: str):
        self.done.add(session_id)
        self.path.write_text(json.dumps(sorted(self.done)))


checkpoint = Checkpoint()

def import_resumable(data: dict) -> dict:
    sid = data["session_id"]
    if checkpoint.is_done(sid):
        logger.info("skip %s (already imported)", sid)
        return {"session_id": sid, "skipped": True}
    result = import_conversation(data)
    if result["failed"] == 0:
        checkpoint.mark(sid)
    return result
```

## Best practices

<AccordionGroup>
  <Accordion title="Sort by timestamp">
    Always sort messages chronologically before importing. Boundary detection
    depends on it. Include real `timestamp`s (in ms) so historical time-based
    recall works.

    ```python theme={null}
    messages.sort(key=lambda m: m["timestamp"])
    ```
  </Accordion>

  <Accordion title="Chunk size">
    Keep chunks at or below the 500-message limit. 100–200 is a good balance of
    throughput and per-request size.
  </Accordion>

  <Accordion title="Concurrency and rate limits">
    Parallelize *across* conversations with the thread pool; keep chunks within a
    conversation ordered. If you hit `429`, lower `max_workers`. The retry
    helper already backs off.
  </Accordion>

  <Accordion title="Stream large directories">
    Read one file at a time (as the importer does) rather than loading every
    export into memory at once.
  </Accordion>
</AccordionGroup>

## Next steps

<CardGroup cols={2}>
  <Card title="Python Integration" icon="python" href="/cookbook/python-integration">
    Client management, error handling, and concurrency patterns.
  </Card>

  <Card title="Multi-Party Conversations" icon="users" href="/cookbook/team-collaboration">
    Work with imported group chat memories.
  </Card>
</CardGroup>


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.