/usr/local/lib64/python3.6/site-packages/caffe2/python
NameSizeModeActions
docs/-0755rm
examples/-0755rm
fakelowp/-0755rm
helpers/-0755rm
ideep/-0755rm
layers/-0755rm
mint/-0755rm
mkl/-0755rm
modeling/-0755rm
models/-0755rm
onnx/-0755rm
operator_test/-0755rm
predictor/-0755rm
rnn/-0755rm
serialized_test/-0755rm
test/-0755rm
trt/-0755rm
__pycache__/-0755rm
allcompare_test.py22550644editdlrm
attention.py123590644editdlrm
benchmark_generator.py49120644editdlrm
binarysize.py55210644editdlrm
brew.py47620644editdlrm
brew_test.py117390644editdlrm
build.py1530644editdlrm
cached_reader.py43940644editdlrm
caffe2_pybind11_state.cpython-36m-x86_64-linux-gnu.so482997120755editdlrm
caffe2_pybind11_state_gpu.cpython-36m-x86_64-linux-gnu.so490481440755editdlrm
caffe_translator.py352270644editdlrm
caffe_translator_test.py35530644editdlrm
checkpoint.py321010644editdlrm
checkpoint_test.py134050644editdlrm
cnn.py76260644editdlrm
context.py28410644editdlrm
context_test.py17920644editdlrm
control.py193090644editdlrm
control_ops_grad.py288930644editdlrm
control_ops_grad_test.py17520644editdlrm
control_ops_util.py108630644editdlrm
control_test.py122760644editdlrm
convert.py550644editdlrm
convert_test.py2010644editdlrm
convnet_benchmarks.py205330644editdlrm
convnet_benchmarks_test.py8390644editdlrm
core.py1194000644editdlrm
core_gradients_test.py380220644editdlrm
core_test.py476830644editdlrm
crf.py132500644editdlrm
crf_predict.py11590644editdlrm
crf_viterbi_test.py16630644editdlrm
dataio.py235320644editdlrm
dataio_test.py175750644editdlrm
dataset.py128860644editdlrm
data_parallel_model.py831000644editdlrm
data_parallel_model_test.py561450644editdlrm
data_workers.py159410644editdlrm
data_workers_test.py65610644editdlrm
db_file_reader.py66080644editdlrm
db_test.py11100644editdlrm
device_checker.py51570644editdlrm
dyndep.py15330644editdlrm
embedding_generation_benchmark.py52560644editdlrm
experiment_util.py36250644editdlrm
extension_loader.py7440644editdlrm
fakefp16_transform_lib.py3220644editdlrm
filler_test.py7480644editdlrm
functional.py44150644editdlrm
functional_test.py42040644editdlrm
fused_8bit_rowwise_conversion_ops_test.py39450644editdlrm
gradient_checker.py153770644editdlrm
gradient_check_test.py207290644editdlrm
gru_cell.py51290644editdlrm
hip_test_util.py4050644editdlrm
hsm_util.py22590644editdlrm
hypothesis_test.py1057620644editdlrm
hypothesis_test_util.py268530644editdlrm
ideep_test_util.py9980644editdlrm
layers_test.py929310644editdlrm
layer_model_helper.py293400644editdlrm
layer_model_instantiator.py39350644editdlrm
layer_parameter_sharing_test.py91480644editdlrm
layer_test_util.py48750644editdlrm
lazy.py2770644editdlrm
lazy_dyndep.py25620644editdlrm
lazy_dyndep_test.py39140644editdlrm
lengths_reducer_fused_8bit_rowwise_ops_test.py75750644editdlrm
lengths_reducer_rowwise_8bit_ops_test.py57100644editdlrm
lstm_benchmark.py106490644editdlrm
memonger.py340410644editdlrm
memonger_test.py369100644editdlrm
mkl_test_util.py11420644editdlrm
model_device_test.py47770644editdlrm
model_helper.py234920644editdlrm
model_helper_test.py23360644editdlrm
modifier_context.py17720644editdlrm
muji.py81310644editdlrm
muji_test.py30580644editdlrm
net_builder.py276790644editdlrm
net_builder_test.py113820644editdlrm
net_drawer.py142640644editdlrm
net_printer.py127040644editdlrm
net_printer_test.py31900644editdlrm
nomnigraph.py42160644editdlrm
nomnigraph_test.py154270644editdlrm
nomnigraph_transformations.py37870644editdlrm
nomnigraph_transformations_test.py57670644editdlrm
normalizer.py14110644editdlrm
normalizer_context.py10070644editdlrm
normalizer_test.py4870644editdlrm
numa_benchmark.py22300644editdlrm
numa_test.py16630644editdlrm
observer_test.py53160644editdlrm
operator_fp_exceptions_test.py12480644editdlrm
optimizer.py788130644editdlrm
optimizer_context.py14620644editdlrm
optimizer_test.py307050644editdlrm
optimizer_test_util.py91870644editdlrm
parallelize_bmuf_distributed_test.py99080644editdlrm
parallel_workers.py76820644editdlrm
parallel_workers_test.py35010644editdlrm
pipeline.py172830644editdlrm
pipeline_test.py25420644editdlrm
predictor_constants.py1980644editdlrm
python_op_test.py91690644editdlrm
queue_util.py44590644editdlrm
record_queue.py44530644editdlrm
recurrent.py132970644editdlrm
regularizer.py211200644editdlrm
regularizer_context.py10130644editdlrm
regularizer_test.py102660644editdlrm
rnn_cell.py682330644editdlrm
schema.py456210644editdlrm
schema_test.py157540644editdlrm
scope.py36230644editdlrm
scope_test.py52490644editdlrm
session.py76420644editdlrm
session_test.py20780644editdlrm
sparse_to_dense_mask_test.py65650644editdlrm
sparse_to_dense_test.py35560644editdlrm
task.py242740644editdlrm
task_test.py8700644editdlrm
test_util.py35240644editdlrm
text_file_reader.py19900644editdlrm
timeout_guard.py40540644editdlrm
toy_regression_test.py28220644editdlrm
transformations.py18320644editdlrm
transformations_test.py119600644editdlrm
tt_core.py93490644editdlrm
tt_core_test.py25180644editdlrm
utils.py141810644editdlrm
utils_test.py13990644editdlrm
visualize.py63150644editdlrm
workspace.py252630644editdlrm
workspace_test.py348440644editdlrm
_import_c_extension.py22500644editdlrm
__init__.py39250644editdlrm
Edit: /usr/local/lib64/python3.6/site-packages/caffe2/python/checkpoint_test.py (13405B)
from caffe2.python.schema import Struct, ConstRecord from caffe2.python import core, workspace, model_helper from caffe2.python.session import LocalSession from caffe2.python.dataset import Dataset from caffe2.python.pipeline import pipe from caffe2.python.checkpoint import ( CheckpointManager, MultiNodeCheckpointManager, Job, JobRunner, epoch_limiter, UploadTaskGroupBuilder, db_name) from caffe2.python.net_builder import ops from caffe2.python.task import Node, Task, TaskGroup, WorkspaceType, Cluster from caffe2.python.test_util import TestCase from caffe2.python.dataio import ReaderWithLimit import numpy as np import os import shutil import tempfile def build_pipeline(node_id): with Node('trainer_%d' % node_id): with Job.current().init_group, Task(): data_arr = Struct(('val', np.array(list(range(10))))) data = ConstRecord(ops, data_arr) ds = Dataset(data, name='dataset:%d' % node_id) full_reader = ds.reader(ops) total = ops.Const([100]) def inc_total(rec): ops.Add([total, rec.val()], [total]) epoch_reader = ReaderWithLimit(full_reader, num_iter=3) pipe(epoch_reader, processor=inc_total) Job.current().add_stop_condition(epoch_reader.data_finished()) return [total] EXPECTED_TOTALS = [103, 115, 136, 145] def local_copy_op(src, dest): def copy_op(inputs, outputs): shutil.copyfile(src, dest) return copy_op class UploadToLocalFile(UploadTaskGroupBuilder): def __init__(self, dest_dir): self.dest_dir = dest_dir def build(self, epoch, checkpoint_manager): with TaskGroup(WorkspaceType.GLOBAL) as upload_task_group: for node, manager in checkpoint_manager._node_managers: with Node(str(node)), Task(): src_path = db_name(epoch, manager._node_name, manager._db_prefix) dest_path = os.path.join(self.dest_dir, str(node)) ops.Python((local_copy_op, [src_path, dest_path], {}))([], []) return upload_task_group class TestCheckpoint(TestCase): def run_with(self, builder): with Cluster(): with Job() as job: outputs = build_pipeline(node_id=0) output_fetcher = Task(step=core.Net('empty'), outputs=outputs) def fetch_total(session): session.run(output_fetcher) return output_fetcher.outputs()[0].fetch() session, checkpoint = builder() job.compile(LocalSession) num_epochs = JobRunner(job, checkpoint).train(session) self.assertEquals(num_epochs, len(EXPECTED_TOTALS)) self.assertEquals(fetch_total(session), EXPECTED_TOTALS[-1]) for initial_epoch in range(1, num_epochs + 1): session, checkpoint = builder() JobRunner( job, checkpoint, resume_from_epoch=initial_epoch ).train(session) self.assertEquals(fetch_total(session), EXPECTED_TOTALS[-1]) for epoch in range(1, num_epochs + 1): session.run(checkpoint.load(epoch)) self.assertEquals(fetch_total(session), EXPECTED_TOTALS[epoch - 1]) def test_single_checkpoint(self): # test single node try: tmpdir = tempfile.mkdtemp() def builder(): ws = workspace.C.Workspace() session = LocalSession(ws) checkpoint = CheckpointManager(tmpdir, 'temp_node', 'minidb') return session, checkpoint self.run_with(builder) finally: shutil.rmtree(tmpdir) # test multi-node try: tmpdir = tempfile.mkdtemp() def builder(): ws = workspace.C.Workspace() session = LocalSession(ws) checkpoint = MultiNodeCheckpointManager(tmpdir, 'minidb') return session, checkpoint self.run_with(builder) finally: shutil.rmtree(tmpdir) def test_ckpt_name_and_load_model_from_ckpts(self): try: num_nodes = 3 tmpdir = tempfile.mkdtemp() # First, check if the checkpoint name generation mechanism is # correct. checkpoint = MultiNodeCheckpointManager(tmpdir, 'minidb') with Cluster(): with Job() as job: for node_id in range(num_nodes): build_pipeline(node_id) job.compile(LocalSession) checkpoint.init(job.nodes_to_checkpoint()) for node_id in range(num_nodes): epoch = 5 node_name = 'trainer_%d' % node_id expected_db_name = tmpdir + '/' + node_name + '.5' self.assertEquals( checkpoint.get_ckpt_db_name(node_name, epoch), expected_db_name) shutil.rmtree(tmpdir) # Next, check mechanism to load model from checkpoints. tmpdir = tempfile.mkdtemp() workspace.ResetWorkspace() for node_id in range(num_nodes): ws = workspace.C.Workspace() session = LocalSession(ws) checkpoint = MultiNodeCheckpointManager(tmpdir, 'minidb') with Cluster(): with Job() as job: build_pipeline(node_id) job.compile(LocalSession) job_runner = JobRunner(job, checkpoint) num_epochs = job_runner.train(session) self.assertEquals(num_epochs, len(EXPECTED_TOTALS)) # There are 17 global blobs after finishing up the job runner. # (only blobs on init_group are checkpointed) self.assertEquals(len(ws.blobs), 17) ws = workspace.C.Workspace() session = LocalSession(ws) self.assertEquals(len(ws.blobs), 0) model_blob_names = ['trainer_1/task_2/GivenTensorInt64Fill:0', 'trainer_2/task_2/GivenTensorInt64Fill:0'] checkpoint = MultiNodeCheckpointManager(tmpdir, 'minidb') with Cluster(): with Job() as job: for node_id in range(num_nodes): build_pipeline(node_id) job.compile(LocalSession) job_runner = JobRunner(job, checkpoint) job_runner.load_blobs_from_checkpoints( blob_names=model_blob_names, epoch=1, session=session) # Check that we can successfully load from checkpoints of epochs # 1 to 4, but not epoch 5. for epoch in range(1, 5): self.assertTrue( job_runner.load_blobs_from_checkpoints( blob_names=model_blob_names, epoch=epoch, session=session)) # Check that all the model blobs are loaded. for blob_name in model_blob_names: self.assertTrue(ws.has_blob(blob_name)) self.assertEquals( ws.fetch_blob(blob_name), np.array([EXPECTED_TOTALS[epoch - 1]])) self.assertFalse( job_runner.load_blobs_from_checkpoints( blob_names=model_blob_names, epoch=5, session=session)) finally: shutil.rmtree(tmpdir) def test_upload_checkpoint(self): try: tmpdir = tempfile.mkdtemp() upload_dir = os.path.join(tmpdir, "upload") os.mkdir(upload_dir) num_nodes = 3 # The uploaded files do not exist yet. for node_id in range(num_nodes): node_name = 'trainer_%d' % node_id upload_path = os.path.join(upload_dir, node_name) self.assertFalse(os.path.exists(upload_path)) # Create and run the job runner. for node_id in range(3): ws = workspace.C.Workspace() session = LocalSession(ws) checkpoint = MultiNodeCheckpointManager(tmpdir, 'minidb') with Cluster(): with Job() as job: build_pipeline(node_id) job.compile(LocalSession) local_upload_builder = UploadToLocalFile(upload_dir) job_runner = JobRunner( job, checkpoint, upload_task_group_builder=local_upload_builder) num_epochs = job_runner.train(session) self.assertEquals(num_epochs, len(EXPECTED_TOTALS)) # The uploaded files should exist now. for node_id in range(num_nodes): node_name = 'trainer_%d' % node_id upload_path = os.path.join(upload_dir, node_name) self.assertTrue(os.path.exists(upload_path)) finally: shutil.rmtree(tmpdir) def test_ckpt_save_failure(self): num_nodes = 3 # The goal of this test is to ensure that the job runs # successfully even if saving a checkpoint fails. # Hence tmpdir is a non existent directory to emulate a failure # while saving checkpoints tmpdir = "/tmp/path_does_not_exist/" # Check the saving checkpoint failure does not cause job failure workspace.ResetWorkspace() for node_id in range(num_nodes): ws = workspace.C.Workspace() session = LocalSession(ws) checkpoint = MultiNodeCheckpointManager(tmpdir, 'minidb') with Cluster(): with Job() as job: build_pipeline(node_id) job.compile(LocalSession) job_runner = JobRunner(job, checkpoint) num_epochs = job_runner.train(session) # make sure all epochs are executed even though saving the checkpoint failed # Saving checkpoint failure should not cause job failure self.assertEquals(num_epochs, len(EXPECTED_TOTALS)) def test_download_group_simple(self): """ A simple test that ensures we have download task group executed between epoch_group and exit_group. """ model = model_helper.ModelHelper(name="test_model") download_net = core.Net("download_net") for name in ["input1", "input2", "output", "download_result"]: model.param_init_net.ConstantFill([], [name], shape=[8, ], value=1.0, run_once=0) model.net.Add(["input1", "input2"], ["output"]) download_net.Copy(["output"], ["download_result"]) # All blob values are initialized as 1.0, after download_net executed # we expect to see download result is the same as training result. with Job() as job: with Node("trainer:0"): with job.init_group: Task(step=model.param_init_net) with job.epoch_group: with Task(): with ops.loop(1): ops.net(model.net) with job.download_group: Task(step=download_net) epoch_limiter(job, 1) ws = workspace.C.Workspace() session = LocalSession(ws) job_runner = JobRunner(job) job_runner.train(session) expected_result = np.full(8, 2.0).astype(np.float32) self.assertTrue(np.array_equal(expected_result, ws.fetch_blob("output"))) self.assertTrue(np.array_equal(expected_result, ws.fetch_blob("download_result"))) def test_reuse_checkpoint_manager(self): """ A simple test that ensures we can reuse a MultiNodeCheckpointManager object. """ try: tmpdir = tempfile.mkdtemp() ws = workspace.C.Workspace() session = LocalSession(ws) checkpoint = MultiNodeCheckpointManager(tmpdir, 'minidb') with Job() as job: outputs = build_pipeline(node_id=0) output_fetcher = Task(step=core.Net('empty'), outputs=outputs) job.compile(LocalSession) def fetch_total(session): session.run(output_fetcher) return output_fetcher.outputs()[0].fetch() num_epochs = JobRunner(job, checkpoint).train(session) for initial_epoch in range(1, num_epochs + 1): JobRunner( job, checkpoint, resume_from_epoch=initial_epoch ).train(session) self.assertEquals(fetch_total(session), EXPECTED_TOTALS[-1]) finally: shutil.rmtree(tmpdir)