Skip to content

Resume a long-running copy from a continuation token.

An operation's monitor URL is the only state needed to keep polling it. A ContinuationToken serializes that URL (with its resolution mode) to JSON, so a job can be picked up by a later call, process or host — the pattern behind durable copy/migration workers. This example submits a copy, persists the token, then deliberately throws the in-memory result away and resumes from the token alone, exactly as a restarted worker would.

Permissions

Requires delegated permission Files.ReadWrite.

Reference

View source

from __future__ import annotations

import argparse
import asyncio
import sys
from pathlib import Path

from office365.graph_client import GraphClient
from office365.runtime.lro import ContinuationToken, OperationPoller, OperationStatus
from tests import create_unique_name
from tests.settings import client_id, password, tenant, username


def report(status: OperationStatus) -> None:
    pct = f"{status.percentage_complete:5.1f}%" if status.percentage_complete is not None else "  ..."
    print(f"\r  {pct}  {status.status or status.http_status}", end="", flush=True)


def write_token(token_path: Path, token: str) -> None:
    token_path.write_text(token)


def read_token(token_path: Path) -> str:
    return token_path.read_text()


def remove_token(token_path: Path) -> None:
    token_path.unlink(missing_ok=True)


async def main() -> None:
    parser = argparse.ArgumentParser(description="Resume a long-running operation from a token")
    parser.add_argument("--token-file", default="lro_continuation.json", help="where to persist the token")
    parser.add_argument("--keep", action="store_true", help="keep the sample files")
    args = parser.parse_args()
    token_path = Path(args.token_file)

    client = GraphClient(tenant=tenant).with_username_and_password(client_id, username, password)
    async with client:
        root = client.me.drive.root
        source = await root.upload(f"{create_unique_name('resume_src')}.txt", b"resume\n").execute_query_async()
        dest = await root.create_folder(create_unique_name("resume_dst")).execute_query_async()
        await dest.get().execute_query_async()

        # -- "Process 1": submit and persist the continuation token, then drop it --
        result = source.copy(name="copy.txt", parent=dest)
        await result.execute_query_async()
        poller = result.to_poller(interval=2, timeout=300)
        write_token(token_path, poller.to_continuation_token().to_json())
        print(f"Submitted copy; token written to {token_path}")

    # -- "Process 2": resume from the token alone, with a fresh client --
    client = GraphClient(tenant=tenant).with_username_and_password(client_id, username, password)
    async with client:
        token = ContinuationToken.from_json(read_token(token_path))
        print(f"Resumed at status={token.status!r}; polling {token.poll_url}")
        poller = OperationPoller.from_continuation_token(client, token, interval=2, timeout=300, authenticate=True)
        status = await poller.wait_async(on_progress=report)
        print(f"\nTerminal: {status.status} (resource_id={status.resource_id})")

        if status.resource_id and not args.keep:
            copied = await client.me.drive.items[status.resource_id].get().execute_query_async()
            await copied.delete_object().execute_query_async()
        if not args.keep:
            source = await client.me.drive.root.get_by_path(source.name).get().execute_query_async()
            await source.delete_object().execute_query_async()
            dest = await client.me.drive.root.get_by_path(dest.name).get().execute_query_async()
            await dest.delete_object().execute_query_async()
            remove_token(token_path)
            print("Samples cleaned up.")


if __name__ == "__main__":
    try:
        asyncio.run(main())
    except KeyboardInterrupt:
        sys.exit(130)

← Back to Async / await