Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
55 changes: 55 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,61 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Features

- Servers now isolate each client's state in a **job**, so a single discipline
server can support concurrent clients. Previously a server held one
discipline instance and shared it: a second client's `Setup` ran
`_clear_data()` and rebuilt `_var_meta` while the first client was still
using it, and `SetOptions` and the in-place shape edits of
`SetVariableShapes` collided the same way. The victim either aborted with
`INTERNAL` on a `KeyError` or, worse, had a short array read into a longer
buffer and received a **zero-padded result with no error** -- which an
optimizer would happily consume. `examples/rosenbrock.py` was the sharp
case, since its variable shape derives from an option, so whichever client
called `SetOptions` last fixed the shapes both clients got.
- A job is a session owning one discipline instance, so the state an author
keeps on `self` -- a mesh, a solver, an `om.Problem` -- is private to one
client. **Every existing discipline hook signature is unchanged**; only the
server constructor moves, from `ExplicitServer(discipline=Paraboloid())` to
`ExplicitServer(discipline=Paraboloid)`. A class works directly as a
factory when its `initialize()` does its own configuration; a discipline
configured externally needs a closure or `functools.partial`. `Discipline`
gains a `job` attribute, giving `self.job.job_id`, and an optional
`teardown_job()` hook called when a job ends or is evicted.
- Three RPCs are added to `DisciplineService`: `StartJob`, `EndJob` and
`KeepAlive`. The job id travels in a `philote-job-id` metadata header rather
than a message field, which costs 7.7 us on a unary call and 13.9 us on a
stream; HPACK indexes the repeated value, so after the first call it is a
byte or two on the wire. A field would instead have been re-serialized on
every chunk of every array with only the first ever read, and each of the
unary RPCs -- all of which take `google.protobuf.Empty` -- would have needed
its own request message. Clients attach the header through a channel
interceptor, so no call site passes it. `GetInfo` and `GetAvailableOptions`
stay job-independent, since they describe the discipline class.
- Clients start a job lazily on the first call that needs one, so existing
scripts and both OpenMDAO components work unchanged. `start_job()`,
`end_job()`, `keep_alive()` and a `job()` context manager are available for
explicit control. An unknown or expired job raises the new `PhiloteJobError`
and is **not** silently replaced: the state that job held is gone, and an
optimizer continuing against a fresh discipline would return plausible but
wrong results. As part of this, every client method now translates gRPC
errors into Philote exceptions; previously only the compute calls did, and a
server-side failure during `run_setup` escaped as a raw `grpc.RpcError`.
- Servers cap concurrent jobs (`max_jobs`, default 8) and evict idle ones
(`ttl`, default one hour), because a job can hold a mesh or a live solver and
a client that dies would otherwise leak it. Exceeding the cap returns
`RESOURCE_EXHAUSTED` rather than exhausting memory. Note that the gRPC
thread pool, not `max_jobs`, is the cap that actually binds -- every in-flight
RPC holds a worker for its whole duration -- so the server warns at startup
when the pool is smaller than the job limit.
- Separate jobs may evaluate concurrently, with no global lock in the path, but
the GIL decides whether that yields throughput. Measured with four
concurrent clients: a pure-Python discipline sees 1.0x (278 ms to 1107 ms), a
NumPy `A @ A` sees 0.8x because threaded BLAS already saturates the cores
with one call, and a discipline whose compiled solver releases the GIL sees
4.0x (306 ms to 310 ms). Jobs buy correctness unconditionally and throughput
conditionally; pure-Python disciplines, including `OpenMdaoSubProblem`, become
correct under concurrent clients rather than faster.

- Continuous array data is now read and written through the packed wire
buffer directly, rather than through the protobuf `repeated double` API,
which converts every element to and from a boxed Python float. A packed
Expand Down
65 changes: 60 additions & 5 deletions docs/docs/getting-started/quickstart.md
Original file line number Diff line number Diff line change
Expand Up @@ -68,20 +68,35 @@ from concurrent import futures
import grpc
# ...

server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
server = grpc.server(futures.ThreadPoolExecutor(max_workers=16))
```

Next, the **Paraboloid** discipline is attached to the server:
Next, the **Paraboloid** discipline is attached to the server. Note the class
`Paraboloid`, not an instance `Paraboloid()` — the server builds one discipline
per job, which is what lets several clients share a server without overwriting
each other's setup, so it needs something it can call:

```python
import philote_mdo.general as pmdo
from philote_mdo.examples import Paraboloid
# ...

discipline = pmdo.ExplicitServer(discipline=Paraboloid())
discipline = pmdo.ExplicitServer(discipline=Paraboloid)
discipline.attach_to_server(server)
```

A class works directly as a factory when its `initialize()` performs its own
configuration. A discipline that has to be configured from the outside needs a
closure or `functools.partial` instead:

```python
from functools import partial

discipline = pmdo.ExplicitServer(
discipline=partial(MyDiscipline, mesh_file="wing.cgns")
)
```

Finally, the port of the server is defined (opening a port is necessary for network communication) and the server is started:

```python
Expand All @@ -106,9 +121,9 @@ import philote_mdo.general as pmdo
from philote_mdo.examples import Paraboloid


server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
server = grpc.server(futures.ThreadPoolExecutor(max_workers=16))

discipline = pmdo.ExplicitServer(discipline=Paraboloid())
discipline = pmdo.ExplicitServer(discipline=Paraboloid)
discipline.attach_to_server(server)

server.add_insecure_port("[::]:50051")
Expand All @@ -117,6 +132,46 @@ print("Server started. Listening on port 50051.")
server.wait_for_termination()
```

## Jobs

Each client gets its own **job** on the server: a session owning one discipline
instance and everything built on it, including the options it set and the
variable metadata `Setup` produced. Two clients can therefore use one server at
the same time without interfering.

The client starts a job on its first call that needs one, so the code above
needs no changes to benefit. When you want to control the lifetime explicitly,
use the context manager, which releases the server's resources on exit:

```python
with client.job():
client.run_setup()
client.get_variable_definitions()
outputs = client.run_compute(inputs)
```

If the server no longer recognises a job -- it expired, or the server was
restarted -- the client raises `PhiloteJobError`. It deliberately does not
start a replacement job on your behalf, because whatever state the old job held
is gone, and an optimizer that carried on would return plausible but wrong
results.

:::note
Servers cap the number of concurrent jobs (`max_jobs`) and evict jobs that go
unused (`ttl`). Make sure the gRPC thread pool has at least as many workers as
the job cap: every in-flight RPC holds a worker for its whole duration, so a
pool smaller than the cap makes jobs queue instead of run. The server warns at
startup when the two are mismatched.
:::

:::warning
Separate jobs may evaluate at the same time, but Python's GIL decides whether
that produces a speedup. A discipline wrapping a compiled solver releases the
GIL and genuinely runs in parallel. A pure-Python discipline does not: jobs make
it *correct* under concurrent clients, not faster. For parallel evaluation of
pure-Python disciplines, run several server processes.
:::

## Calling the Discipline Using a Client

Now that a server is running, it can be queried using a client. Philote-Python offers a number of clients for this purpose, ranging from the general implementation to OpenMDAO and CSDL components. However, under the hood, the OpenMDAO and CSDL components use the general client implementation.
Expand Down
15 changes: 7 additions & 8 deletions docs/docs/tutorials/implicit-disciplines.md
Original file line number Diff line number Diff line change
Expand Up @@ -227,14 +227,13 @@ import grpc
import philote_mdo.general as pmdo

def run_server():
# Create the discipline
discipline = QuadraticSolver()

# Create gRPC server
server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
# Create gRPC server. The pool should be at least as large as the
# server's job cap, since every in-flight RPC holds a worker.
server = grpc.server(futures.ThreadPoolExecutor(max_workers=16))

# Create and attach implicit server
impl_server = pmdo.ImplicitServer(discipline=discipline)
# Create and attach implicit server. It takes a factory and builds one
# discipline per job, so concurrent clients do not interfere.
impl_server = pmdo.ImplicitServer(discipline=QuadraticSolver)
impl_server.attach_to_server(server)

# Start server
Expand Down Expand Up @@ -268,7 +267,7 @@ def run_production_server():
]
)

impl_server = pmdo.ImplicitServer(discipline=discipline)
impl_server = pmdo.ImplicitServer(discipline=QuadraticSolver)
impl_server.attach_to_server(server)

# Use secure connection in production
Expand Down
10 changes: 8 additions & 2 deletions examples/openmdao_sellar.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,9 +34,15 @@


def run():
server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
# the pool must be at least as large as the server's job cap: every
# in-flight RPC holds a worker for its whole duration
server = grpc.server(futures.ThreadPoolExecutor(max_workers=16))

discipline = pmdo.ExplicitServer(discipline=SellarGroup())
# One discipline instance is built per job, so concurrent clients do not
# interfere. A class works as the factory here because its initialize()
# does its own configuration; a discipline configured from the outside
# needs a closure or functools.partial.
discipline = pmdo.ExplicitServer(discipline=SellarGroup)
discipline.attach_to_server(server)

server.add_insecure_port("[::]:50051")
Expand Down
10 changes: 8 additions & 2 deletions examples/parabaloid_explicit.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,9 +34,15 @@


def run():
server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
# the pool must be at least as large as the server's job cap: every
# in-flight RPC holds a worker for its whole duration
server = grpc.server(futures.ThreadPoolExecutor(max_workers=16))

discipline = pmdo.ExplicitServer(discipline=Paraboloid())
# One discipline instance is built per job, so concurrent clients do not
# interfere. A class works as the factory here because its initialize()
# does its own configuration; a discipline configured from the outside
# needs a closure or functools.partial.
discipline = pmdo.ExplicitServer(discipline=Paraboloid)
discipline.attach_to_server(server)

server.add_insecure_port("[::]:50051")
Expand Down
10 changes: 8 additions & 2 deletions examples/quadratic_implicit.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,8 +34,14 @@


def run():
server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
discipline = pmdo.ImplicitServer(discipline=QuadradicImplicit())
# the pool must be at least as large as the server's job cap: every
# in-flight RPC holds a worker for its whole duration
server = grpc.server(futures.ThreadPoolExecutor(max_workers=16))
# One discipline instance is built per job, so concurrent clients do not
# interfere. A class works as the factory here because its initialize()
# does its own configuration; a discipline configured from the outside
# needs a closure or functools.partial.
discipline = pmdo.ImplicitServer(discipline=QuadradicImplicit)
discipline.attach_to_server(server)

server.add_insecure_port("[::]:50051")
Expand Down
2 changes: 2 additions & 0 deletions philote_mdo/general/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,8 @@
from .explicit_server import ExplicitServer
from .implicit_server import ImplicitServer

from .job import JOB_METADATA_KEY, Job, JobState, JobStore

from .discipline import Discipline
from .explicit_discipline import ExplicitDiscipline
from .implicit_discipline import ImplicitDiscipline
16 changes: 16 additions & 0 deletions philote_mdo/general/discipline.py
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,11 @@ def __init__(self):
# flag that indicates the discipline is implicit
self._is_implicit = False

# the job that owns this instance, assigned by the server when the
# job is created. One discipline instance serves exactly one job, so
# anything stored on self is private to that client.
self.job = None

# dictionary of available discipline options (with types)
self.options_list = {}

Expand Down Expand Up @@ -283,6 +288,17 @@ def setup_partials(self):
def configure(self):
pass

def teardown_job(self):
"""
Releases whatever this instance holds, before its job is discarded.

Called when the client ends the job and when the server evicts one
that has gone idle, so a discipline that opened a file, started a
subprocess, or built a solver should close it here. Overriding this
is optional; the default does nothing.
"""
pass

def _clear_data(self):
"""
Clears all metadata of the discipline.
Expand Down
Loading
Loading