Skip to content

Latest commit

 

History

History
429 lines (306 loc) · 14 KB

File metadata and controls

429 lines (306 loc) · 14 KB

App API Setup Guide

The App API packages the current Python project, registers a temporary Flame application from the flmrun template, and exposes Python functions and classes as remote services.

Prerequisites

  • A running Flame cluster with the session manager, executor manager, and object cache.
  • A configured flmrun application registered from <config-dir>/applications/flmrun.yaml when the session manager starts.
  • A Python environment that can import flamepy.

Verify the template application:

import flamepy

flmrun = flamepy.get_application("flmrun")
if flmrun is None:
    raise RuntimeError("flmrun application is not registered")

Client Configuration

App reads ~/.flame/flame.yaml through flamepy.core.FlameContext. A minimal local configuration is:

---
current-context: flame
contexts:
  - name: flame
    cluster:
      endpoint: "http://127.0.0.1:8080"
    cache:
      endpoint: "grpc://127.0.0.1:9090"

package.storage is optional. When it is absent, App uploads packages to the Flame object cache through cache.endpoint.

Supported package storage schemes:

  • grpc:// and grpcs://: Flame object cache storage.
  • file://: shared filesystem storage.
  • http:// and https://: HTTP storage with PUT, GET, and DELETE support.

Example with explicit HTTP package storage:

package:
  storage: "http://127.0.0.1:5050/packages"

The optional app field defaults to flmrun. Set it only when the cluster uses a custom template application name:

app: "custom-flmrun"

Basic Usage

Verify A Cluster

The Python SDK includes a small App E2E command for checking that a configured cluster can run App workloads:

uv run --project e2e python -m e2e.app

The check packages a temporary project, verifies the configured flmrun template, runs a function service, passes ObjectFuture values between services, and runs a remotely constructed class service. Use --json for machine-readable output, --tasks N to change the function-task count, and --python-version 3.12 to verify a specific executor Python runtime.

Function Service

import flamepy.app as app

app.init("sum-app")


@app.service()
def sum_fn(a: int, b: int) -> int:
    return a + b


print(sum_fn(1, 3))  # local call
result = sum_fn.remote(1, 3)
print(result.get())
app.destroy()

Functions default to autoscaling sessions.

Class Service

import flamepy.app as app


app.init("counter-app")


@app.service(warmup=1)
class Counter:
    def __init__(self, initial: int = 0):
        self._count = initial

    def add(self, value: int) -> int:
        self._count += value
        return self._count

    def get(self) -> int:
        return self._count

counter = Counter.remote(10)
counter.add(1).wait()
counter.add(3).wait()
counter.add(5).wait()
print(counter.get().get())
app.destroy()

Decorating Counter declares the service class but does not create a session. Counter(10) constructs a local object. Counter.remote(10) creates the remote proxy and session, and flmrun runs Counter.__init__(10) in the executor. Call methods directly on the proxy; class-level calls such as Counter.add(...) are not supported. The options on @app.service(...) apply to the session created for the proxy.

Passing ObjectFuture Values

Remote calls return ObjectFuture. App sends a pickled ServiceRequest and ServiceResponse through the task RPC. Each successful response contains a Ref: ValueRef for an inline value or an explicit ObjectRef returned by the worker. Both reference types can be passed as arguments and resolve through flamepy.core.get_object() or flamepy.core.aio.get_object().

A pickled request over 128 KiB raises before task submission. A pickled response over 128 KiB fails the remote task and raises when its future is read with get(), ref(), or wait(). Worker exceptions follow the same task failure path; they are not wrapped in ServiceResponse. Cache large values explicitly with flamepy.core.put_object and pass the ObjectRef as an argument.

update_object, patch_object, and download_object require an ObjectRef; a ValueRef has no stored cache object to modify or download.

import flamepy.app as app


app.init("chain-app")


@app.service()
def double(value: int) -> int:
    return value * 2


@app.service()
def add(a: int, b: int) -> int:
    return a + b


first = double.remote(21)
total = add.remote(first, 8)
print(total.get())
app.destroy()

Services can also call one another from an executor. A captured service proxy is serialized as a reference to the already-created session:

@app.service()
def fn_a(value: int) -> int:
    return value * 2


fn_a_service = app.remote(fn_a)


@app.service()
def fn_b(value: int) -> int:
    return fn_a_service(value).get()


print(fn_b.remote(21).get())

API Reference

Process-wide application

Initialize the application once before declaring services:

flamepy.app.init(name, fail_if_exists=False)
  • name: application and package name. Repeating init() with the same name returns the same runtime handle.
  • fail_if_exists: when True, raise if the application already exists. The default reuses an existing enabled application and skips cleanup for it; a disabled application always produces an error.

Top-level functions:

  • app.service(autoscale=None, warmup=0, resreq=None): canonical decorator. Decoration declares a function or class service without creating a session. Call the decorated object directly for local execution. For functions, fn.remote(*args, **kwargs) submits a remote call on a shared session. app.remote(fn) creates an independent function proxy. For classes, Clazz.remote(*args, **kwargs) or app.remote(Clazz, *args, **kwargs) creates a proxy and passes constructor arguments to the executor.
  • get(futures): resolve multiple ObjectFuture values to concrete objects.
  • ref(futures): return Ref values, either inline ValueRef or cache-backed ObjectRef.
  • wait(futures): wait for multiple futures without fetching objects.
  • select(futures): iterate over futures as they complete.
  • flamepy.app.put(obj): store a shared object under the active application prefix.
  • destroy(): close services, unregister the application, delete cached objects, and remove the uploaded package.

Recursive services are the lifecycle exception. A function call through .remote(...) on a declaration made while service code is running borrows that invocation's existing session. A nested class .remote(...) call creates a borrowed proxy. These calls require neither app.init() nor app.destroy() and do not own or close the parent session. Nested declarations reuse the existing session configuration, so they do not accept autoscale, warmup, or resreq. They must remain in the invocation's execution context, and nested calls must complete before the parent invocation returns.

resreq accepts a resource string, for example:

import flamepy.app as app

app.init("cpu-app")


@app.service(resreq="cpu=1,mem=1g")
def add(left: int, right: int) -> int:
    return left + right

ServiceInstance

app.service() is the canonical decorator. Calling a decorated function or class directly executes it locally. fn.remote(*args, **kwargs) submits a remote function call on a shared session and returns ObjectFuture. app.remote(fn) creates an independent ServiceInstance that can be called repeatedly. Clazz.remote(*args, **kwargs) and app.remote(Clazz, *args, **kwargs) each create a ServiceInstance and send constructor arguments to the executor. Each remote class proxy has its own service ID and retained object; creating two proxies from the same decorated class does not share object state. app.destroy() closes those service sessions.

Public methods on a decorated class must not use names reserved by ServiceInstance, such as close. A collision is rejected when the service is declared rather than silently hiding the user method.

When declarations live in another module, initialize the application before importing that module so its decorators execute against the active application:

import flamepy.app as app

app.init("pipeline-app")

# pipeline.services uses @app.service() for its declarations.
from pipeline import services  # noqa: E402

result = services.transform.remote("input")
print(result.get())
app.destroy()
  • Decorated functions and classes are callable locally.
  • Function proxies from app.remote(fn) and class proxies from either remote factory run calls remotely.
  • Class proxies expose one wrapper method for each public method.
  • Every remote call returns ObjectFuture.

Default service behavior with warmup=0:

Execution object Default autoscale Effective min_instances Effective max_instances
Function or builtin True 0 unlimited
Class True 0 unlimited

For functions, builtins, and classes, autoscale is configurable. When warmup=N and N > 0, autoscaled services use min_instances=N and no max limit; fixed services use min_instances=N and max_instances=N. Passing an already constructed object to app.service() is not supported.

Data-Aware Instance Selection

App services can be plain functions or classes. During a service invocation, call app.publish_attributes(attrs) to add opaque bytes keys held by that executor. Repeated calls during one method accumulate, and App drains the published set after each response. Session Manager unions every response into the executor's retained attribute set. app.session_context() exposes the active Flame session context during the same invocation:

import os
from pathlib import Path

import flamepy.app as app
from flamepy import TaskOptions


app.init("cache-app")


class CacheBase:
    def __init__(self):
        # App retains class execution objects with the executor process.
        self.path = Path(f"/tmp/flame-app-cache-{os.getpid()}")
        self.keys = set()

    def store(self, key: bytes, value: str) -> str:
        self.path.write_text(value)
        self.keys = {key}
        app.publish_attributes(self.keys)
        return value

    def load(self, key: bytes) -> str:
        self.keys.add(key)
        app.publish_attributes(self.keys)
        return self.path.read_text()


@app.service(warmup=2)
class WarmCache(CacheBase):
    pass


@app.service(warmup=0)
class TargetCache(CacheBase):
    pass


key = b"model:block:7"
warm_cache = WarmCache.remote()
target_cache = TargetCache.remote()
warm_cache.store(key, "value").wait()
warm_cache.close()
result = target_cache.load(key, option=TaskOptions(affinity={key}))
print(result.get())
app.destroy()

App publishes accumulated keys after every task invocation, including a failed task. Session Manager retains previously accepted keys, and publishing nothing is a no-op. Keys must be nonempty and at most 256 bytes. One publication round supports at most 1,024 distinct keys and 64 KiB after deduplication.

App constructs class execution objects from the supplied constructor arguments in each executor process. It retains an object by its handle's service ID when that executor binds to another session for the same service. Distinct handles have distinct service IDs and do not share objects, even if they were created from the same decorated class with identical arguments. Functions remain session-scoped and can use their module-level state. The second session is created before the first one closes so its task is ready to use the retained executor. Idle instances remain reusable for twice the application's configured delay_release before Shuffle releases them: 120 seconds with the default 60-second setting.

ObjectFuture

Methods:

  • get(): fetch and deserialize the concrete object.
  • ref(): return a flamepy.core.Ref: ValueRef for inline results or the explicit ObjectRef returned by the worker.
  • wait(): wait for completion without fetching the object.

Package Contents

App packages the current working directory into dist/<name>.tar.gz.

Default exclusions include:

  • .venv, venv
  • __pycache__
  • .pytest_cache, .ruff_cache, .mypy_cache
  • *.egg-info
  • .git, .tox
  • node_modules
  • *.pyc, *.pyo
  • .DS_Store

Add .flmignore or .flameignore to the project directory for more exclusions. Both files use Gitignore patterns, including ! negation and rules in nested directories. For example:

*.log
!important.log
*.pkl
data/

Place the file in the project directory that app.init() packages. When both names exist in the same directory, .flameignore takes precedence. The same files control flmctl deploy --application <directory>; .gitignore and global Git exclusions do not affect either archive.

Working Directory

App derives the registered application's working_directory from the flmrun template. If the template has a working directory, App appends /<app-name>. If the template has no working directory, App leaves it unset.

Troubleshooting

Failed to get application template 'flmrun': confirm the session manager is running and flmrun appears in flmctl list -a.

Storage not configured: configure cache.endpoint or package.storage. In a local setup, cache.endpoint: "grpc://127.0.0.1:9090" is enough.

Storage directory does not exist: for file:// storage, create the directory on a filesystem that both the client and executor nodes can access.

Package upload or download failures: verify the selected storage backend is reachable from both the client and executor nodes. For object-cache storage, check the object-cache service on port 9090.

Pickle or import errors: keep service functions and classes importable from the packaged project, and make sure executor nodes have the required Python dependencies.

See Also