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
8 changes: 7 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -69,4 +69,10 @@ trace-*
test_sherlock/

# SMS API #
.hpc_env
.hpc_env

# Local analysis notebooks
*.ipynb

# Generated simulation state/output
data/*.json
26 changes: 26 additions & 0 deletions configs/N_all_media_conditions.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
{
"experiment_id": "all_media_conditions2",
"sim_data_path": null,
"suffix_time": false,
"lineage_seed": 0,
"generations": 8,
"n_init_sims": 1,
"single_daughters": true,
"emitter": "parquet",
"emitter_arg": {
"out_dir": "out"
},
"variants": {
"condition": {
"condition": {
"value": ["with_aa", "acetate", "succinate", "no_oxygen"]
}
}
},
"analysis_options": {
"multivariant": {
"doubling_time_hist": {"skip_n_gens": 0},
"doubling_time_line": {}
}
}
}
29 changes: 29 additions & 0 deletions configs/colony_baseline_test.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
{
"inherit_from": ["spatial.json"],
"description": "Test: 2 generations (4 cells), save parquet per cell, save final state",

"seed": 0,
"sim_data_path": "/home/katha/dev/vEcoli-colony-sim/out/all_media_conditions2/parca/kb/simData.cPickle",

"emitter": "parquet",
"emitter_arg": {
"out_dir": "out/colony_runs/baseline_test1"
},

"inner_emitter": "parquet",

"max_duration": 6000,

"save": true,
"save_times": [6000],
"colony_save_prefix": "baseline_test1",

"parallel": false,

"engine_process_reports": [
["boundary"],
["bulk"],
["listeners"],
["environment", "exchange"]
]
}
30 changes: 30 additions & 0 deletions configs/colony_baseline_test2.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
{
"inherit_from": ["spatial.json"],
"description": "Test2: read from trial1, run for another generation, save parquet per cell, save final state",
"initial_colony_file": "baseline_2gen_seed_0_colony_t6000",

"seed": 0,
"sim_data_path": "/home/katha/dev/vEcoli-colony-sim/out/all_media_conditions2/parca/kb/simData.cPickle",

"emitter": "parquet",
"emit_step": 60,
"emitter_arg": {
"out_dir": "out/colony_runs/baseline_3rd_gen_seed_0"
},
"emit_config": false,

"max_duration": 9000,

"save": true,
"save_times": [9000],
"colony_save_prefix": "baseline_3rd_gen",

"parallel": false,

"engine_process_reports": [
["boundary"],
["bulk"],
["listeners"],
["environment", "exchange"]
]
}
30 changes: 30 additions & 0 deletions configs/colony_baseline_test3.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
{
"inherit_from": ["spatial.json"],
"description": "Test3: 2 generations (4 cells), save parquet per cell, save final state",

"seed": 10,
"sim_data_path": "/home/katha/dev/vEcoli-colony-sim/out/all_media_conditions2/parca/kb/simData.cPickle",

"emit_step": 60,
"emitter": "parquet",
"emitter_arg": {
"out_dir": "out/colony_runs/baseline_test3"
},

"inner_emitter": "parquet",

"max_duration": 6000,

"save": true,
"save_times": [6000],
"colony_save_prefix": "baseline_test3",

"parallel": false,

"engine_process_reports": [
["boundary"],
["bulk"],
["listeners"],
["environment", "exchange"]
]
}
30 changes: 30 additions & 0 deletions configs/colony_baseline_test4.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
{
"inherit_from": ["spatial.json"],
"description": "Test4: Emit_step test for only 120 seconds, save parquet per cell, save final state",

"seed": 10,
"sim_data_path": "/home/katha/dev/vEcoli-colony-sim/out/all_media_conditions2/parca/kb/simData.cPickle",

"emit_step": 60,
"emitter": "parquet",
"emitter_arg": {
"out_dir": "out/colony_runs/baseline_test4_2"
},

"inner_emitter": "parquet",

"max_duration": 120,

"save": true,
"save_times": [120],
"colony_save_prefix": "baseline_test4",

"parallel": false,

"engine_process_reports": [
["boundary"],
["bulk"],
["listeners"],
["environment", "exchange"]
]
}
71 changes: 61 additions & 10 deletions ecoli/experiments/ecoli_engine_process.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,9 @@
tuplify_topology,
)
from ecoli.library.logging_tools import write_json

# 27 August Fix3: needed so shutdown can finalize parquet emitters.
from ecoli.library.parquet_emitter import ParquetEmitter
from ecoli.library.sim_data import RAND_MAX
from ecoli.library.schema import not_a_process
from ecoli.library.json_state import get_state_from_file
Expand Down Expand Up @@ -149,6 +152,8 @@ def generate_processes(self, config):
def generate_topology(self, config):
pass

# Fix 3: pass the emit cadence through to the inner EngineProcess.


class EcoliEngineProcess(Composer):
"""
Expand Down Expand Up @@ -202,6 +207,8 @@ def generate_processes(self, config):
"tunnel_out_schemas": config["tunnel_out_schemas"],
"stub_schemas": config["stub_schemas"],
"seed": (config["seed"] + 1) % RAND_MAX,
# Fix 3: pass the emit cadence through to the inner EngineProcess.
"emit_step": config.get("emit_step", 1),
"inner_emitter": config["inner_emitter"],
"divide": config["divide"],
# Inner sim will set store at ('division_trigger',)
Expand Down Expand Up @@ -263,30 +270,54 @@ def colony_save_states(engine, config):
engine.update(time_to_next_save)
time_elapsed = config["save_times"][i]

# Save the full state of the super-simulation
# Save the full state of the super-simulation # Fix 1
state_to_save = engine.state.get_value(condition=not_a_process)

# Get internal state from the EngineProcess sub-simulation
for agent_id in state_to_save["agents"]:
engine.state.get_path(
("agents", agent_id, "cell_process")
).value.send_command("get_inner_state")
for agent_id in state_to_save["agents"]:
# The parquet emitter wraps the entire outer simulation in
# agents/outer. Colony save files should NOT contain this wrapper,
# because initial_colony_file expects agents/<agent_id>.
if config["emitter"] == "parquet":
colony_state = state_to_save["agents"]["outer"]
cell_process_base_path = ("agents", "outer", "agents")
else:
colony_state = state_to_save
cell_process_base_path = ("agents",)

# Get internal state from each EngineProcess sub-simulation
for agent_id in colony_state["agents"]:
cell_process_path = cell_process_base_path + (
agent_id,
"cell_process",
)
engine.state.get_path(cell_process_path).value.send_command(
"get_inner_state"
)

for agent_id in colony_state["agents"]:
cell_process_path = cell_process_base_path + (
agent_id,
"cell_process",
)
cell_state = engine.state.get_path(
("agents", agent_id, "cell_process")
cell_process_path
).value.get_command_result()

# Can't save, but will be restored when loading state
del cell_state["environment"]["exchange_data"]

# Shared processes are re-initialized on load
del cell_state["process"]

# Save bulk and unique dtypes
cell_state["bulk_dtypes"] = str(cell_state["bulk"].dtype)
cell_state["unique_dtypes"] = {}
for name, mols in cell_state["unique"].items():
cell_state["unique_dtypes"][name] = str(mols.dtype)
state_to_save["agents"][agent_id] = cell_state

state_to_save = serialize_value(state_to_save)
colony_state["agents"][agent_id] = cell_state

state_to_save = serialize_value(colony_state) # Fix 1

if config.get("colony_save_prefix", None):
write_json(
"data/"
Expand Down Expand Up @@ -317,6 +348,19 @@ def colony_save_states(engine, config):
engine.update(time_remaining)


def finalize_parquet_emitters(processes, seen=None): # Fix 3
if seen is None:
seen = set()
if isinstance(processes, dict):
for process in processes.values():
finalize_parquet_emitters(process, seen)
elif hasattr(processes, "emitter"):
emitter = processes.emitter
if isinstance(emitter, ParquetEmitter) and id(emitter) not in seen:
seen.add(id(emitter))
emitter.finalize()


def run_simulation(config):
"""
Main method for running colony simulations in
Expand Down Expand Up @@ -368,6 +412,8 @@ def run_simulation(config):
"stub_schemas": stub_schemas,
"parallel": config["parallel"],
"divide": config["divide"],
# Fix 3: pass the emit cadence into the inner EngineProcess setup.
"emit_step": config.get("emit_step", 1),
"tunnels_in": (
("environment",),
("boundary",),
Expand Down Expand Up @@ -512,6 +558,11 @@ def run_simulation(config):
else:
engine.update(config["max_duration"])
engine.end()
# Fix 2&3: flush buffered parquet output on normal shutdown.
if config["emitter"] == "parquet":
engine.emitter.success = True
engine.emitter.finalize()
finalize_parquet_emitters(engine.processes)

if config["profile"]:
report_profiling(engine.stats)
Expand Down
22 changes: 12 additions & 10 deletions ecoli/processes/engine_process.py
Original file line number Diff line number Diff line change
Expand Up @@ -210,6 +210,7 @@ class EngineProcess(Process):
"inner_composer_config": {},
"outer_composer_config": {},
"seed": 0,
"emit_step": 1, # Fix3
"inner_emitter": "null",
"divide": False,
"division_threshold": None,
Expand Down Expand Up @@ -514,16 +515,17 @@ def next_update(self, timestep, states):
# other words, since we rely on the outer Engine to apply the
# updates, we have to wait for those updates from the previous
# timestep to be applied before we emit data.
data = self.sim.state.emit_data()
if isinstance(self.emitter, ParquetEmitter):
# ParquetEmitter expects per-agent data under agents/<agent_id>.
data = {"agents": {self.parameters["agent_id"]: data}}
data["time"] = self.sim.global_time
emit_config = {
"table": "history",
"data": data,
}
self.emitter.emit(emit_config)
if self.sim.global_time % self.parameters["emit_step"] == 0: # Fix 3
data = self.sim.state.emit_data()
if isinstance(self.emitter, ParquetEmitter):
# ParquetEmitter expects per-agent data under agents/<agent_id>.
data = {"agents": {self.parameters["agent_id"]: data}}
data["time"] = self.sim.global_time
emit_config = {
"table": "history",
"data": data,
}
self.emitter.emit(emit_config)

# Run inner simulation for timestep.
try:
Expand Down
14 changes: 14 additions & 0 deletions reconstruction/ecoli/dataclasses/process/two_component_system.py
Original file line number Diff line number Diff line change
Expand Up @@ -248,6 +248,20 @@ def __setstate__(self, state):
self._populate_derivative_and_jacobian()
self.dependency_matrix = self._make_dependency_matrix()

@property
def modified_molecules(
self,
): # Fix 0
value = self.__dict__.get("modified_molecules")
if value is None:
value = self.make_modified_molecule_list()
self.__dict__["modified_molecules"] = value
return value

@modified_molecules.setter # Fix 0
def modified_molecules(self, value):
self.__dict__["modified_molecules"] = value

def _buildComplexToMonomer(self, sim_data, modifiedFormsMonomers, tcsMolecules):
"""
Maps each complex to a dictionary that maps each subunit of the complex
Expand Down
Loading