Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
264 commits
Select commit Hold shift + click to select a range
7f589f2
vine_schedule_count_committable_cores
Jun 20, 2025
84c135c
Merge branch 'cooperative-computing-lab:master' into separate_pending…
JinZhou5042 Jun 27, 2025
f944fc6
Merge branch 'cooperative-computing-lab:master' into separate_pending…
JinZhou5042 Jun 27, 2025
a217aa8
merge
Jun 30, 2025
ed57fe0
Merge remote-tracking branch 'origin/separate_pending_and_ready_tasks…
Jun 30, 2025
8b9c37c
Merge branch 'cooperative-computing-lab:master' into separate_pending…
JinZhou5042 Jul 5, 2025
c9ccf6a
Merge branch 'cooperative-computing-lab:master' into separate_pending…
JinZhou5042 Sep 23, 2025
6714bca
vine: function infile load mode
Oct 13, 2025
a78dc95
specify utf-8
Oct 13, 2025
83cd050
utf-8
Oct 13, 2025
d6f2fea
vine: some bug fixes for the library code
Oct 13, 2025
d79b1ed
dttools: progress_bar src and test
Oct 14, 2025
354a6d4
add unit
Oct 14, 2025
81b228f
format issue
Oct 14, 2025
d5dad06
vine: valid link when sending message on worker
Oct 14, 2025
5b5c142
merge
Oct 14, 2025
a599721
Merge remote-tracking branch 'origin/progress-bar' into task-graph
Oct 14, 2025
ad3f0fb
Merge remote-tracking branch 'origin/worker-bug-fix' into task-graph
Oct 14, 2025
39b1d66
vine: LIST_ITERATE_REVERSE
Oct 14, 2025
e94875e
temp
Oct 14, 2025
15fe7fd
Merge remote-tracking branch 'origin/LIST_ITERATE_REVERSE' into task-…
Oct 14, 2025
40e9679
vine: task graph executor in C
Oct 14, 2025
21b25a4
simple fix
Oct 14, 2025
648071f
braces in block
Oct 14, 2025
8242136
type fix
Oct 14, 2025
c841f6c
type fix
Oct 14, 2025
39168ef
lint
Oct 14, 2025
e8c23e6
some fixes
Oct 15, 2025
ef0e8a2
bug
Oct 15, 2025
6b527b1
use proxy function and library
Oct 15, 2025
b4e5a7b
comment functions
Oct 15, 2025
1955fd1
use a separate directory
Oct 15, 2025
6b1a5db
backup
Oct 16, 2025
8cd855f
libtask
Oct 16, 2025
5cb9b49
dagvine
Oct 16, 2025
38c7cff
new interface
Oct 16, 2025
dbe57a8
api change
Oct 16, 2025
6bf5de3
new
Oct 16, 2025
2032815
vine task graph tuning
Oct 16, 2025
642d144
new version
Oct 17, 2025
3fab110
remove _task_graph
Oct 18, 2025
5b85f46
new
Oct 18, 2025
ec96041
reg
Oct 18, 2025
767988f
sog && reg
Oct 18, 2025
fa50958
sog
Oct 18, 2025
0a0dd40
milestone
Oct 19, 2025
73308fb
new files
Oct 19, 2025
90476b8
lint and format
Oct 19, 2025
256c02f
rename executor
Oct 19, 2025
ea9285f
checkpoint works
Oct 19, 2025
35243c5
update apis
Oct 19, 2025
7ba4fac
make
Oct 19, 2025
728531b
lint
Oct 19, 2025
ac40b74
make format
Oct 19, 2025
76d55c8
comment
Oct 20, 2025
31af425
Merge branch 'cooperative-computing-lab:master' into task-graph
JinZhou5042 Oct 20, 2025
945e26f
file name fix
Oct 20, 2025
5d00612
Merge branch 'cooperative-computing-lab:master' into separate_pending…
JinZhou5042 Oct 21, 2025
c668184
rename to vinedag
Oct 21, 2025
b8bca0d
add three fields
Oct 21, 2025
4beda4d
Merge remote-tracking branch 'origin/separate_pending_and_ready_tasks…
Oct 21, 2025
1575045
calculate time spent on scheduling
Oct 21, 2025
eabf0c3
Merge remote-tracking branch 'origin/separate_pending_and_ready_tasks…
Oct 21, 2025
c107a9c
gpus
Oct 22, 2025
ec0c8c5
Merge remote-tracking branch 'origin/separate_pending_and_ready_tasks…
Oct 22, 2025
cc27109
time metrics
Oct 29, 2025
a9b95b2
optimized put
Oct 29, 2025
363b5e7
lint
Oct 29, 2025
c7e9d65
lint
Oct 29, 2025
2faa0d1
time metrics
Oct 29, 2025
fd06a26
new structure
Nov 1, 2025
cca3170
git ignore
Nov 1, 2025
b8a3472
Merge branch 'cooperative-computing-lab:master' into task-graph
JinZhou5042 Nov 3, 2025
c1b4e67
revert vine_manager_put_task
Nov 3, 2025
51c13bc
Merge remote-tracking branch 'origin/task-graph' into task-graph
Nov 3, 2025
c0cc116
traverse w->current_libraries instead of w->current_tasks
Nov 3, 2025
2430286
Merge remote-tracking branch 'origin/use-current-library' into task-g…
Nov 3, 2025
a8dd136
merge
Nov 3, 2025
0372e72
revert push_task_to_ready_tasks
Nov 3, 2025
1059321
vine: multiply task priority on resource exhaustion
Nov 3, 2025
409fdd9
Merge remote-tracking branch 'origin/push_task_to_ready_tasks' into t…
Nov 3, 2025
7850141
vine: clean redundant input replicas upon task completion
Nov 3, 2025
3bcc9d0
lint
Nov 3, 2025
128afce
Merge branch 'cooperative-computing-lab:master' into clean_redundant_…
JinZhou5042 Nov 4, 2025
203736a
make it an argument
Nov 4, 2025
af1eb9d
shift-disk-load
Nov 4, 2025
8aab2df
merge
Nov 4, 2025
e3e51b5
lint
Nov 4, 2025
f6e2b0e
comment
Nov 4, 2025
caff977
vine_temp
Nov 4, 2025
cdd96a8
vine: vine_temp.c
Nov 4, 2025
f590c77
merge
Nov 4, 2025
d10c7ac
lint
Nov 4, 2025
57673e7
trigger rebuild
Nov 4, 2025
1d91659
comments
Nov 5, 2025
ed1fc23
remove unrelated code
Nov 5, 2025
17da0b5
Merge remote-tracking branch 'origin/master' into progress-bar
Nov 5, 2025
8280ef6
fix flicker
Nov 5, 2025
ddc460c
update doc
Nov 6, 2025
176201f
fix comment
Nov 6, 2025
b6880b2
trigger rebuild
Nov 6, 2025
062c4d3
Merge branch 'cooperative-computing-lab:master' into library-network-…
JinZhou5042 Nov 10, 2025
ba0b985
review fix
Nov 10, 2025
289e8ee
merge
Dec 3, 2025
c89bd85
modify project name
Dec 3, 2025
111f989
Merge branch 'cooperative-computing-lab:master' into library-network-…
JinZhou5042 Dec 4, 2025
86cbc13
fix
Dec 4, 2025
67b29a5
fix comments
Jan 5, 2026
876a24e
Merge branch 'cooperative-computing-lab:master' into vine_temp
JinZhou5042 Jan 5, 2026
40a0929
merge
Jan 5, 2026
b1e70df
clean
Jan 5, 2026
581208e
merge
Jan 5, 2026
56892f2
merge
Jan 5, 2026
b5d8a0f
merge
Jan 5, 2026
f953106
Merge remote-tracking branch 'origin/library-network-code-bugfix' int…
Jan 5, 2026
058466e
update
Jan 5, 2026
da8aa18
rename recovery_tasks_submitted
Jan 5, 2026
39ae787
timestamp format
Jan 5, 2026
1b0cdf1
vine: check whether hash table key is present before inserting
Jan 6, 2026
ff7eeb2
fix hash table insertion
Jan 6, 2026
91ed9df
Merge remote-tracking branch 'origin/duplicate-key-insertion' into ta…
Jan 6, 2026
0dfdbed
Merge branch 'cooperative-computing-lab:master' into library-infile-l…
JinZhou5042 Jan 6, 2026
933946b
merge
Jan 6, 2026
fe47036
Add example target to Makefile for running example_blueprint
Jan 6, 2026
bc442cd
Add example target to Makefile for running example_blueprint
Jan 6, 2026
747a647
temp
Jan 7, 2026
ffda413
merge
Jan 7, 2026
97993b4
support new dask expr
Jan 8, 2026
2a6d2ca
revert vine_worker
Jan 8, 2026
7d729ab
support dask expressions
Jan 8, 2026
a7d05b1
gitignore
Jan 8, 2026
23331bf
remove context graph from makefile
Jan 8, 2026
6f6aca5
convert bg to dask expr
Jan 9, 2026
a4d078d
support internal in/out
Jan 9, 2026
160c2b9
auto recovery
Jan 9, 2026
f43bd83
rename to User
Jan 9, 2026
845e0a6
block rec tasks w missing inputs
Jan 9, 2026
0810e46
merge
Jan 9, 2026
8e4e54d
Merge branch 'library-infile-load-mode' of github.com:JinZhou5042/cct…
Jan 9, 2026
cad7e29
rename infile_path
Jan 9, 2026
3036a0d
Merge remote-tracking branch 'origin/master' into task-graph
Jan 12, 2026
3abe87b
add compute makespan
Jan 15, 2026
4222532
use processing_transfers in vine_cache_check_files
Jan 15, 2026
4857b76
Merge branch 'cooperative-computing-lab:master' into task-graph
JinZhou5042 Feb 16, 2026
1a937cd
fast task committment
Nov 3, 2025
47904db
checkpoint
Feb 16, 2026
1b9663d
fix var name
Feb 19, 2026
227b25f
comment out time metrics
Feb 19, 2026
8398137
Merge branch 'cooperative-computing-lab:master' into task-graph
JinZhou5042 Mar 23, 2026
b1475cc
ckpt
Mar 19, 2026
07712ac
checkpoint: good condition
Apr 23, 2026
be4d2b1
apply gitignore
Apr 23, 2026
474367f
add graph
Apr 23, 2026
d3a6aee
reapply .gitignore
Apr 24, 2026
171fbed
rm ignore .so files
Apr 24, 2026
ea48d03
Merge branch 'cooperative-computing-lab:master' into master
JinZhou5042 Apr 24, 2026
60821e1
rearrange makefiles
Apr 29, 2026
59ee50d
rename files and dirs
Apr 29, 2026
6804546
rename dirs
Apr 29, 2026
0b7f83e
clean proxy and adaptor
Apr 29, 2026
721bd66
add utils
Apr 29, 2026
602421f
revert vine_manager_put
Apr 29, 2026
b93d400
add dask executor
May 2, 2026
6c5f4fd
clean unused include
May 2, 2026
2d837a2
rm some unused funcs
May 2, 2026
32bfe0d
rm time_end_execution
May 4, 2026
97d9eb5
Merge branch 'cooperative-computing-lab:master' into task-graph
JinZhou5042 May 4, 2026
b82997d
apply gitignore
Apr 23, 2026
0903c50
add graph
Apr 23, 2026
bac5db4
reapply .gitignore
Apr 24, 2026
6d25742
rm ignore .so files
Apr 24, 2026
e35c44d
add back transfer_port_bind_failed
May 4, 2026
6694544
rm unused vars
May 4, 2026
ef29f7c
rm unsed vars
May 4, 2026
db6b758
rm vine_schedule_count_committable_cores
May 4, 2026
b381359
no use cvine at dagvine.py
May 4, 2026
9667cb2
revert some files
May 4, 2026
9dd4df6
rm unsed vars
May 4, 2026
d6b7e34
delete context graph
May 4, 2026
d0c2e51
revert daskvine
May 4, 2026
42344a1
rm binary file
May 4, 2026
e91b470
fix cvine path
May 4, 2026
7afd2ec
clean some includes
May 4, 2026
4006f10
clean comments
May 4, 2026
c60137c
upgrade runner task
May 5, 2026
d02a71d
runtime task creation and io mounting
May 5, 2026
2936fa1
profiling pre/post processing
May 6, 2026
15f9a25
modulize pre/post
May 6, 2026
aff3195
remove comment
May 6, 2026
88ba3ec
chain grouping
May 6, 2026
aa95c24
Merge branch 'cooperative-computing-lab:master' into master
JinZhou5042 May 8, 2026
8f086ae
checkpoint-v0
Jun 8, 2026
84426dc
Merge branch 'cooperative-computing-lab:master' into master
JinZhou5042 Jun 17, 2026
a7758e7
Merge branch 'cooperative-computing-lab:master' into task-graph
JinZhou5042 Jun 17, 2026
ca2959f
renaming files and dirs and function names stuff
Jun 17, 2026
b376824
move project path
Jun 17, 2026
ccbfd42
apply gitignore
Apr 23, 2026
8bb46b2
add graph
Apr 23, 2026
a8fcdbe
reapply .gitignore
Apr 24, 2026
5a41f73
rm ignore .so files
Apr 24, 2026
e7fac33
rm graph/
Jun 17, 2026
bc32231
move to python bindings
Jun 17, 2026
e2a6ffb
add some patches
Jun 23, 2026
0964f61
Merge branch 'cooperative-computing-lab:master' into master
JinZhou5042 Jun 23, 2026
d139e0b
merge
Jun 23, 2026
d6b8145
merge
Jun 23, 2026
12fcd4b
rm stale files
Jun 23, 2026
a026bf0
remove a couple useless code
Jun 23, 2026
accf9e4
no enable debug log
Jun 23, 2026
61fd96b
vine: return recovery tasks
Jun 23, 2026
84acc64
lint
Jun 23, 2026
ebc4a2b
makefile patch
Jun 24, 2026
b315519
remove unused code
Jun 24, 2026
e735769
rename variable name
Jun 24, 2026
7ef6353
Trigger CI rebuild
Jun 24, 2026
4267276
merge
Jun 25, 2026
3a1b98d
Merge remote-tracking branch 'origin/master' into task-graph
Jun 25, 2026
b657c86
Merge remote-tracking branch 'origin/library-infile-load-mode' into t…
Jun 25, 2026
1eb3833
rm stale api
Jun 25, 2026
5dedfd6
vine_cache_check_xfer
Jun 25, 2026
c7ed9f5
add a comment
Jun 25, 2026
5055135
rename to external_recovery_handling
Jun 25, 2026
d1af1c8
Merge branch 'cooperative-computing-lab:master' into return-recovery-…
JinZhou5042 Jun 25, 2026
a12c297
new changes
Jun 25, 2026
947350b
merge
Jun 25, 2026
ae5bed2
reset task commit
Jun 25, 2026
fea3204
rename
Jul 14, 2026
a19c194
ckpt
Jul 16, 2026
345dc7f
merge
Jul 29, 2026
39354a4
Merge branch 'cooperative-computing-lab:master' into task-graph
JinZhou5042 Aug 27, 2026
1b1b399
Merge branch 'cooperative-computing-lab:master' into task-graph
JinZhou5042 Sep 1, 2026
2960cc8
Merge remote-tracking branch 'origin/master' into task-graph
Sep 1, 2026
6a8ee02
vine: defer failed library cleanup during iteration
Sep 1, 2026
dbe6c7b
Merge branch 'library-cleanup' into task-graph
Sep 1, 2026
141a823
vine graph: hoist workflow module dependencies
Sep 1, 2026
8f5d6de
vine graph: stage extension for development builds
Sep 1, 2026
5dab71e
docs: add tested Vine Graph user manual
Sep 1, 2026
303a888
docs: tighten Vine Graph manual prose
Sep 1, 2026
6899979
docs: remove source-tree setup from Vine Graph manual
Sep 1, 2026
b4654d3
docs: remove Vine Graph troubleshooting section
Sep 2, 2026
aabebca
docs: omit task grouping from Vine Graph manual
Sep 2, 2026
2018fb6
Merge remote-tracking branch 'origin/task-graph' into task-graph
Sep 2, 2026
7109c36
docs: simplify Vine Graph user examples
Sep 2, 2026
cfc090c
docs: add remote worker setup to Vine Graph manual
Sep 2, 2026
a177fb4
docs: use manager names in Vine Graph worker example
Sep 2, 2026
3548ef3
docs: make worker shutdown user controlled
Sep 2, 2026
191031b
docs: align Vine Graph worker core counts
Sep 2, 2026
c06fc2e
docs: trim Vine Graph argument internals
Sep 2, 2026
79140b7
docs: remove Vine Graph test summary
Sep 2, 2026
e9deeb4
docs: omit Vine Graph repeat controls
Sep 2, 2026
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
296 changes: 296 additions & 0 deletions doc/manuals/taskvine/vine-graph.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,296 @@
# Vine Graph

Vine Graph describes a directed acyclic graph (DAG) of Python function calls and
runs it with TaskVine. A `Workflow` records the tasks and their dependencies;
`VineGraph` runs the graph and returns the requested results.

## The programming model

The main objects are:

- `Workflow`: owns a graph and all handles created for that graph.
- `TaskHandle`: identifies a task. `workflow.add_task()` returns one.
- `TaskOutputHandle`: represents a task's Python return value, or a selected
part of that value. Obtain one with `task.output()`.
- `FileHandle`: represents an existing frontend file or a file produced in a
task sandbox.
- `VineGraph`: a TaskVine manager that executes a completed workflow.

`VineGraph.run()` is synchronous: it returns after execution finishes. The
returned dictionary is keyed by the target handles supplied by the caller.

## First workflow: local execution

Local execution runs the graph in the manager process and does not need a
worker. It is handy for checking graph construction and task functions. Create
`vine_graph_local.py`:

```python
from ndcctools.taskvine.vine_graph import VineGraph, Workflow


def make_record(value):
return {"values": [value, value + 1], "metadata": {"count": 2}}


def scaled_sum(values, count, scale=1):
assert len(values) == count
return sum(values) * scale


workflow = Workflow()
record = workflow.add_task(make_record, 20)
answer = workflow.add_task(
scaled_sum,
record.output()["values"],
record.output()["metadata"]["count"],
scale=2,
)

with VineGraph(port=0) as manager:
results = manager.run(
workflow,
targets=[answer],
params={
"local-execute": 1,
"output-dir": "./vine-graph-local-output",
},
)

assert results[answer] == 82
print(results[answer])
```

Run it with:

```bash
python vine_graph_local.py
```

Passing `record` directly would be an error because a `TaskHandle` identifies a
task, not its value. Pass `record.output()` instead. Indexing an output handle,
as in `record.output()["values"]`, selects a value inside the consuming task
without adding another graph node.

To request every terminal result, pass `targets=workflow.sink_tasks()`.
`sink_tasks()` returns handles for all tasks with no downstream consumers.

## Passing files between tasks

Use `workflow.file(path)` for an existing frontend file. Use
`task.file(relative_path)` for a file that the task will create inside its
sandbox. A consumer receives either kind of handle as a local path string.

For example, `vine_graph_files.py` passes a frontend file to one task and its
output file to another:

```python
from pathlib import Path

from ndcctools.taskvine.vine_graph import VineGraph, Workflow


def uppercase(source_path):
text = Path(source_path).read_text().upper()
Path("result.txt").write_text(text)
return len(text)


def verify(result_path, expected_length):
text = Path(result_path).read_text()
assert len(text) == expected_length
return text


Path("input.txt").write_text("vine graph\n")

workflow = Workflow()
source = workflow.file("input.txt")
producer = workflow.add_task(uppercase, source)
produced_file = producer.file("result.txt")
consumer = workflow.add_task(verify, produced_file, producer.output())

with VineGraph(port=0) as manager:
results = manager.run(
workflow,
targets=[consumer, produced_file],
params={
"local-execute": 1,
"output-dir": "./vine-graph-file-output",
},
)

assert results[consumer] == "VINE GRAPH\n"
assert Path(results[produced_file]).read_text() == "VINE GRAPH\n"
print(results[consumer], end="")
```

Task output paths must be non-empty relative paths that stay inside the task
sandbox. Each output path may be declared only once for a given task. Handles
belong to one `Workflow` and cannot be passed into another workflow.

## Distributed execution with a worker

Distributed mode is the default. The manager creates a TaskVine task-runner
library, and workers execute the graph nodes. Create
`vine_graph_distributed.py`:

```python
from ndcctools.taskvine.vine_graph import VineGraph, Workflow


def square(value):
return value * value


def add_all(*values):
return sum(values)


workflow = Workflow()
squares = [workflow.add_task(square, value) for value in range(1, 6)]
total = workflow.add_task(add_all, *(task.output() for task in squares))

with VineGraph(port=0, name="vine-graph-example") as manager:
results = manager.run(
workflow,
targets=[total],
params={
"libcores": 4,
"output-dir": "./vine-graph-distributed-output",
},
)

assert results[total] == 55
print(results[total])
```

Run the manager in one terminal:

```bash
python vine_graph_distributed.py
```

Start a worker in another terminal:

```bash
vine_worker -M vine-graph-example --cores 4
```

The manager name is used by remote workers to find the manager through the
TaskVine catalog. Stop the worker with `Ctrl-C`.

`--cores` on `vine_worker`, `vine_submit_workers`, and `vine_factory` must match
the manager's `libcores`. All examples here use 4.

## HTCondor workers

Start the manager:

```bash
python vine_graph_distributed.py
```

In another shell, submit workers to HTCondor:

```bash
vine_submit_workers -T condor -M vine-graph-example \
--cores 4 5
```

This submits five workers with four cores each. Use `condor_q` to check their
status. The submit command prints the cluster ID; use `condor_rm CLUSTER_ID` to
stop the workers.

## Using a factory

A factory adds and removes workers as the workload changes. Keep the manager
running, then start the factory in another shell:

```bash
vine_factory -T condor -M vine-graph-example \
--min-workers 1 --max-workers 10 --cores 4
```

To use the same conda environment on the workers, make a Poncho tarball:

```bash
poncho_package_create --ignore-editable-packages \
"$CONDA_PREFIX" vine-graph-env.tar.gz
```

Pass it to the factory:

```bash
vine_factory -T condor -M vine-graph-example \
--min-workers 1 --max-workers 10 --cores 4 \
--poncho-env vine-graph-env.tar.gz
```

Stop the factory with `Ctrl-C` when the workflow is finished.

## Dask graphs

Set `from_dask=True` to convert a Dask-style graph. Low-level task dictionaries
and common Dask collection forms are supported. For example:

```python
from ndcctools.taskvine.vine_graph import VineGraph


def increment(value):
return value + 1


def multiply(left, right):
return left * right


dask_graph = {
"incremented": (increment, 1),
"answer": (multiply, "incremented", 10),
}

with VineGraph(port=0) as manager:
results = manager.run(
dask_graph,
targets=["answer"],
from_dask=True,
params={
"local-execute": 1,
"output-dir": "./vine-graph-dask-output",
},
)

assert results["answer"] == 20
print(results["answer"])
```

For Dask collections, pass a dictionary of Dask collection objects. Do not mix
collection and non-collection values in that dictionary.

## Execution parameters

Parameters may be supplied through the `params` argument to `run()` or through
`manager.set_params()` before `run()`.

Useful parameters include:

| Parameter | Purpose |
| --- | --- |
| `local-execute` | Run in process when set to `1`; use TaskVine workers when `0`. |
| `output-dir` | Store serialized results; local mode also stores task sandboxes here. |
| `checkpoint-dir` | Store executor checkpoints here. |
| `libcores` | Number of cores assigned to the task-runner library. |

## Run the project regression tests

From the repository root:

```bash
cd taskvine/test
./TR_vine_graph_workflow_examples.sh prepare
./TR_vine_graph_workflow_examples.sh run
./TR_vine_graph_dask_adaptor.sh prepare
./TR_vine_graph_dask_adaptor.sh run
```
30 changes: 24 additions & 6 deletions poncho/src/poncho/library_network_code.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,9 @@
r, w = os.pipe()
exec_method = None

# infile load mode for function tasks inside this library
function_infile_load_mode = None


# This class captures how results from FunctionCalls are conveyed from
# the library to the manager.
Expand Down Expand Up @@ -85,6 +88,18 @@ def sigchld_handler(signum, frame):
os.write(w, b"a")


# Load the infile for a function task inside this library
def load_function_infile(in_file_path):
if function_infile_load_mode == "cloudpickle":
with open(in_file_path, "rb") as f:
return cloudpickle.load(f)
elif function_infile_load_mode == "json":
with open(in_file_path, "r", encoding="utf-8") as f:
return json.load(f)
else:
raise ValueError(f"invalid infile load mode: {function_infile_load_mode}")


# Read data from worker, start function, and dump result to `outfile`.
def start_function(in_pipe_fd, thread_limit=1):
# read length of buffer to read
Expand Down Expand Up @@ -131,8 +146,7 @@ def start_function(in_pipe_fd, thread_limit=1):
os.chdir(function_sandbox)

# parameters are represented as infile.
with open("infile", "rb") as f:
event = cloudpickle.load(f)
event = load_function_infile("infile")

# output of execution should be dumped to outfile.
result = globals()[function_name](event)
Expand All @@ -158,11 +172,10 @@ def start_function(in_pipe_fd, thread_limit=1):
return -1, function_id
elif exec_method == "fork":
try:
arg_infile = os.path.join(function_sandbox, "infile")
with open(arg_infile, "rb") as f:
event = cloudpickle.load(f)
infile_path = os.path.join(function_sandbox, "infile")
event = load_function_infile(infile_path)
except Exception:
stdout_timed_message(f"TASK {function_id} error: can't load the arguments from {arg_infile}")
stdout_timed_message(f"TASK {function_id} error: can't load the arguments from {infile_path}")
return -1, function_id
p = os.fork()
if p == 0:
Expand Down Expand Up @@ -368,11 +381,16 @@ def main():
global exec_method
exec_method = library_info['exec_mode']

# set infile load mode of functions in this library
global function_infile_load_mode
function_infile_load_mode = library_info['function_infile_load_mode']

# send configuration of library, just its name for now
config = {
"name": library_info['library_name'],
"taskid": args.task_id,
"exec_mode": exec_method,
"function_infile_load_mode": function_infile_load_mode,
}
send_configuration(config, out_pipe_fd, args.worker_pid)

Expand Down
Loading
Loading