Files
revng-revng/tests/pypeline/test_all.py
2025-09-10 12:05:15 +02:00

649 lines
23 KiB
Python

#
# This file is distributed under the MIT License. See LICENSE.md for details.
#
# pylint: disable=redefined-outer-name
from __future__ import annotations
import os
from tempfile import NamedTemporaryFile
from typing import Optional, TypeVar, Union
import pytest
from simple_pipeline import ChildDictContainer, DictModel, GeneratorPipe, InPlacePipe, MyKind
from simple_pipeline import MyObjectID, NullAnalysis, PurgeAllAnalysis, PurgeOneAnalysis
from simple_pipeline import RootDictContainer, SameKindPipe, ToHigherKindPipe, ToLowerKindPipe
from revng.pypeline import initialize_pypeline
from revng.pypeline.analysis import AnalysisBinding
from revng.pypeline.container import ContainerDeclaration
from revng.pypeline.model import ReadOnlyModel
from revng.pypeline.object import Kind, ObjectSet
from revng.pypeline.pipeline import Artifact, Pipeline
from revng.pypeline.pipeline_node import PipelineConfiguration, PipelineNode
from revng.pypeline.pipeline_parser import load_pipeline_yaml_file
from revng.pypeline.storage.memory import InMemoryStorageProvider
from revng.pypeline.storage.sqlite3 import SQlite3StorageProvider
from revng.pypeline.storage.storage_provider import ContainerLocation, SavePointsRange
from revng.pypeline.storage.storage_provider import StorageProvider
from revng.pypeline.task.pipe import Pipe
from revng.pypeline.task.requests import Requests
from revng.pypeline.task.savepoint import SavePoint
# Fill the registries
initialize_pypeline()
Value = Union[str, int]
T = TypeVar("T")
def mandatory(arg: Optional[T]) -> T:
assert arg is not None
return arg
@pytest.fixture
def model():
return DictModel()
@pytest.fixture(params=["memory", "sqlite3"])
def storage_provider(request):
storage_provider: StorageProvider
if request.param == "memory":
storage_provider = InMemoryStorageProvider()
yield storage_provider
elif request.param == "sqlite3":
with NamedTemporaryFile() as f:
storage_provider = SQlite3StorageProvider(":memory:", f.name)
yield storage_provider
else:
raise ValueError()
def test_kind():
# Test rank
assert MyKind.ROOT.rank() == 0
assert MyKind.CHILD.rank() == 1
assert MyKind.GRANDCHILD.rank() == 2
assert MyKind.CHILD2.rank() == 1
assert MyKind.root() == MyKind.ROOT
assert MyKind.CHILD.parent() == MyKind.ROOT
assert MyKind.GRANDCHILD.parent() == MyKind.CHILD
assert MyKind.CHILD2.parent() == MyKind.ROOT
assert MyKind.kinds() == [MyKind.ROOT, MyKind.CHILD, MyKind.GRANDCHILD, MyKind.CHILD2]
# Test relation
assert MyKind.ROOT.relation(MyKind.ROOT)[0] == Kind.Relation.SAME
assert MyKind.CHILD.relation(MyKind.CHILD2)[0] == Kind.Relation.UNRELATED
assert MyKind.ROOT.relation(MyKind.CHILD) == (
Kind.Relation.ANCESTOR,
[MyKind.ROOT, MyKind.CHILD],
)
assert MyKind.CHILD.relation(MyKind.ROOT) == (
Kind.Relation.DESCENDANT,
[MyKind.CHILD, MyKind.ROOT],
)
assert MyKind.ROOT.relation(MyKind.GRANDCHILD) == (
Kind.Relation.ANCESTOR,
[MyKind.ROOT, MyKind.CHILD, MyKind.GRANDCHILD],
)
assert MyKind.GRANDCHILD.relation(MyKind.ROOT) == (
Kind.Relation.DESCENDANT,
[MyKind.GRANDCHILD, MyKind.CHILD, MyKind.ROOT],
)
def test_pipe_prerequisites_for(model) -> None:
# These are special declarations and have to exactly match the args of the
# pipes being tested.
root = ObjectSet(MyKind.ROOT, {MyObjectID.root()})
one_two = ObjectSet(
MyKind.CHILD, {MyObjectID(MyKind.CHILD, "one"), MyObjectID(MyKind.CHILD, "two")}
)
one_two_three = ObjectSet(
MyKind.CHILD,
{
MyObjectID(MyKind.CHILD, "one"),
MyObjectID(MyKind.CHILD, "two"),
MyObjectID(MyKind.CHILD, "three"),
},
)
empty_child_set = ObjectSet(MyKind.CHILD, set())
empty_root_set = ObjectSet(MyKind.ROOT, set())
pipe: Pipe
pipe = InPlacePipe()
result = pipe.prerequisites_for(model, [one_two])
assert result == [one_two]
pipe = SameKindPipe()
result = pipe.prerequisites_for(model, [empty_child_set, one_two])
assert result == [one_two, empty_child_set]
pipe = ToLowerKindPipe()
result = pipe.prerequisites_for(model, [empty_child_set, root])
assert result == [one_two_three, empty_root_set]
pipe = ToHigherKindPipe()
result = pipe.prerequisites_for(model, [empty_root_set, one_two_three])
assert result == [root, empty_child_set]
def test_savepoint_prerequisites_for(storage_provider, model) -> None:
storage_provider.set_model(model.serialize())
child: ContainerDeclaration = ContainerDeclaration("child", ChildDictContainer)
save_point = SavePoint("save", [child])
container = ChildDictContainer()
container.add_object(MyObjectID(MyKind.CHILD, "one"))
requests = Requests({child: ObjectSet(MyKind.CHILD, {MyObjectID(MyKind.CHILD, "one")})})
savepoint_range = SavePointsRange(start=0, end=0)
configuration_id = "132124ujh12jk4jk124kj"
# The storage is empty so the savepoint should be transparent
result = save_point.prerequisites_for(
requests=requests,
configuration_id=configuration_id,
storage_provider=storage_provider,
savepoint_range=savepoint_range,
)
assert result == requests
# Force the savepoint to store the object
save_point.run(
containers={child: container},
incoming=requests,
outgoing=requests,
configuration_id=configuration_id,
storage_provider=storage_provider,
savepoint_range=savepoint_range,
)
result = save_point.prerequisites_for(
requests=requests,
configuration_id=configuration_id,
storage_provider=storage_provider,
savepoint_range=savepoint_range,
)
assert result == Requests({child: ObjectSet(MyKind.CHILD, set())})
result = save_point.prerequisites_for(
requests=Requests(
{
child: ObjectSet(
MyKind.CHILD, {MyObjectID(MyKind.CHILD, "one"), MyObjectID(MyKind.CHILD, "two")}
)
}
),
configuration_id=configuration_id,
storage_provider=storage_provider,
savepoint_range=savepoint_range,
)
expected = Requests({child: ObjectSet(MyKind.CHILD, {MyObjectID(MyKind.CHILD, "two")})})
assert result == expected
def test_pipeline_inplace(model, storage_provider):
child_cont: ContainerDeclaration = ContainerDeclaration("arg", ChildDictContainer)
declarations = [child_cont]
pipeline_configuration: PipelineConfiguration = {}
storage_provider.set_model(model.serialize())
one = ObjectSet(MyKind.CHILD, {MyObjectID(MyKind.CHILD, "one")})
one_two = ObjectSet(
MyKind.CHILD, {MyObjectID(MyKind.CHILD, "one"), MyObjectID(MyKind.CHILD, "two")}
)
# Create the pipeline
begin_node = PipelineNode(SavePoint("begin", to_save=declarations))
inplace_node = PipelineNode(InPlacePipe(), bindings=[child_cont])
end_node = PipelineNode(SavePoint("end", to_save=declarations))
begin_node.add_successor(inplace_node).add_successor(end_node)
pipeline: Pipeline = Pipeline(set(declarations), begin_node)
assert begin_node.savepoint_range == SavePointsRange(start=1, end=2)
assert inplace_node.savepoint_range == SavePointsRange(start=1, end=2)
assert end_node.savepoint_range == SavePointsRange(start=2, end=2)
# Force the savepoint to store the objects
container = ChildDictContainer()
container.add_object(MyObjectID(MyKind.CHILD, "one"))
container.add_object(MyObjectID(MyKind.CHILD, "two"))
begin_configuration_id = begin_node.configuration_id(pipeline_configuration)
begin_node.run(
model=ReadOnlyModel(model),
containers={child_cont: container},
incoming=Requests({child_cont: one_two}),
outgoing=Requests({child_cont: one_two}),
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
assert set(
storage_provider.has(
location=ContainerLocation(
savepoint_id=1,
container_id="arg",
configuration_id=begin_configuration_id,
),
keys=one_two,
)
) == set(one_two)
containers = pipeline.schedule(
model=ReadOnlyModel(model),
target_node=end_node,
requests=Requests({child_cont: one}),
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
).run(
model=ReadOnlyModel(model),
storage_provider=storage_provider,
)
assert containers[child_cont].objects() == one
end_configuration_id = end_node.configuration_id(pipeline_configuration)
assert list(
storage_provider.has(
location=ContainerLocation(
savepoint_id=2,
container_id="arg",
configuration_id=end_configuration_id,
),
keys=one,
)
) == list(one)
def test_pipeline_up_down(model, storage_provider):
root1: ContainerDeclaration = ContainerDeclaration("root_source", RootDictContainer)
root2: ContainerDeclaration = ContainerDeclaration("root_destination", RootDictContainer)
child1: ContainerDeclaration = ContainerDeclaration("child_destination", ChildDictContainer)
child2: ContainerDeclaration = ContainerDeclaration("child_source", ChildDictContainer)
declarations = [root1, root2, child1, child2]
pipeline_configuration: PipelineConfiguration = {}
storage_provider.set_model(model.serialize())
begin_node = PipelineNode(SavePoint("begin", declarations))
up_node = PipelineNode(ToHigherKindPipe(), bindings=[root1, child1])
same_node = PipelineNode(SameKindPipe(), bindings=[child1, child2])
down_node = PipelineNode(ToLowerKindPipe(), bindings=[child2, root2])
end_node = PipelineNode(SavePoint("end", declarations))
begin_node.add_successor(up_node).add_successor(same_node).add_successor(
down_node
).add_successor(end_node)
pipeline: Pipeline = Pipeline(set(declarations), begin_node)
assert begin_node.savepoint_range == SavePointsRange(start=1, end=2)
assert up_node.savepoint_range == SavePointsRange(start=1, end=2)
assert same_node.savepoint_range == SavePointsRange(start=1, end=2)
assert down_node.savepoint_range == SavePointsRange(start=1, end=2)
assert end_node.savepoint_range == SavePointsRange(start=2, end=2)
root_obj = ObjectSet(MyKind.ROOT, {MyObjectID.root()})
# Force the savepoint to store the objects
container = RootDictContainer()
container.add_object(MyObjectID.root())
requests = Requests({root1: root_obj})
begin_node.run(
model=ReadOnlyModel(model),
containers={root1: container},
incoming=requests,
outgoing=requests,
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
containers = pipeline.schedule(
model=ReadOnlyModel(model),
target_node=end_node,
requests=Requests({root2: root_obj}),
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
).run(
model=ReadOnlyModel(model),
storage_provider=storage_provider,
)
assert containers[root2].objects() == root_obj
end_configuration_id = end_node.configuration_id(pipeline_configuration)
assert set(
storage_provider.has(
location=ContainerLocation(
savepoint_id=2,
container_id="root_destination",
configuration_id=end_configuration_id,
),
keys=[MyObjectID.root()],
)
) == {MyObjectID.root()}
def test_artifact(model, storage_provider):
child1: ContainerDeclaration = ContainerDeclaration("source", ChildDictContainer)
child2: ContainerDeclaration = ContainerDeclaration("destination", ChildDictContainer)
declarations = [child1, child2]
pipeline_configuration: PipelineConfiguration = {}
one = ObjectSet.from_list([MyObjectID(MyKind.CHILD, "one")])
storage_provider.set_model(model.serialize())
begin_node = PipelineNode(SavePoint("begin", declarations))
same_node = PipelineNode(SameKindPipe(), bindings=[child1, child2])
begin_node.add_successor(same_node)
artifact = Artifact("artifact", same_node, child2)
pipeline: Pipeline = Pipeline(set(declarations), begin_node)
assert begin_node.savepoint_range == SavePointsRange(start=1, end=1)
assert same_node.savepoint_range == SavePointsRange(start=1, end=1)
# Force the savepoint to store the objects
container = ChildDictContainer()
container.add_object(MyObjectID(MyKind.CHILD, "one"))
requests = Requests({child1: one})
begin_node.run(
model=ReadOnlyModel(model),
containers={child1: container},
incoming=requests,
outgoing=requests,
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
res = pipeline.get_artifact(
model=ReadOnlyModel(model),
artifact=artifact,
requests=one,
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
assert res.objects() == one
def test_invalidation(model, storage_provider):
# TODO: simple test: one pipe, one save point, multiple objects, invalidate one
child: ContainerDeclaration = ContainerDeclaration("arg", ChildDictContainer)
declarations = [child]
pipeline_configuration: PipelineConfiguration = {}
storage_provider.set_model(model.serialize())
object_one = MyObjectID(MyKind.CHILD, "one")
expected_output = ObjectSet(MyKind.CHILD, {object_one})
pipe = PipelineNode(GeneratorPipe(), bindings=[child])
savepoint = PipelineNode(SavePoint("end", declarations))
pipe.add_successor(savepoint)
pipeline: Pipeline = Pipeline(set(declarations), pipe)
assert pipe.savepoint_range == SavePointsRange(start=0, end=1)
assert savepoint.savepoint_range == SavePointsRange(start=1, end=1)
containers = pipeline.schedule(
model=ReadOnlyModel(model),
target_node=savepoint,
requests=Requests({child: expected_output}),
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
).run(
model=ReadOnlyModel(model),
storage_provider=storage_provider,
)
assert containers[child].objects() == expected_output
assert set(
storage_provider.has(
location=ContainerLocation(
savepoint_id=1,
container_id="arg",
configuration_id=savepoint.configuration_id(pipeline_configuration),
),
keys=expected_output,
)
) == set(expected_output)
storage_provider.invalidate({"/two"})
assert set(
storage_provider.has(
location=ContainerLocation(
savepoint_id=1,
container_id="arg",
configuration_id=savepoint.configuration_id(pipeline_configuration),
),
keys=expected_output,
)
) == set(expected_output)
# Change something we depend on
storage_provider.invalidate({"/one"})
assert not list(
storage_provider.has(
location=ContainerLocation(
savepoint_id=1,
container_id="arg",
configuration_id=savepoint.configuration_id(pipeline_configuration),
),
keys=expected_output,
)
)
def test_analysis(model, storage_provider):
# TODO: simple test: one pipe, one save point, multiple objects, invalidate one
child: ContainerDeclaration = ContainerDeclaration("arg", ChildDictContainer)
declarations = [child]
pipeline_configuration: PipelineConfiguration = {}
model["/test/test"] = "test"
model["/test/test2"] = "test2"
storage_provider.set_model(model.serialize())
object_one = MyObjectID(MyKind.CHILD, "one")
expected_output = ObjectSet(MyKind.CHILD, {object_one})
pipe = PipelineNode(
GeneratorPipe(),
bindings=[child],
)
savepoint = PipelineNode(SavePoint("end", declarations))
pipe.add_successor(savepoint)
pipeline: Pipeline = Pipeline(
set(declarations),
pipe,
analyses={
AnalysisBinding(
NullAnalysis(name="null_analysis"),
(child,),
savepoint,
),
AnalysisBinding(
PurgeAllAnalysis(name="purge_all_analysis"),
(child,),
savepoint,
),
AnalysisBinding(
PurgeOneAnalysis(name="purge_one_analysis"),
(child,),
savepoint,
),
},
)
assert pipe.savepoint_range == SavePointsRange(start=0, end=1)
assert savepoint.savepoint_range == SavePointsRange(start=1, end=1)
orig_model = model.clone()
new_model = pipeline.run_analysis(
model=ReadOnlyModel(model),
analysis_name="null_analysis",
requests=Requests({child: expected_output}),
analysis_configuration="",
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
assert model == orig_model, "Model should not be modified by the analysis"
assert new_model == orig_model, "This analysis doesn't invalidate anything"
new_model = pipeline.run_analysis(
model=ReadOnlyModel(model),
analysis_name="purge_all_analysis",
requests=Requests({child: expected_output}),
analysis_configuration="",
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
assert model == orig_model, "Model should not be modified by the analysis"
assert new_model == DictModel(), "This analysis invalidates everything"
new_model = pipeline.run_analysis(
model=ReadOnlyModel(model),
analysis_name="purge_one_analysis",
requests=Requests({child: expected_output}),
analysis_configuration="",
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
assert model == orig_model, "Model should not be modified by the analysis"
expected_model = DictModel()
expected_model["/test/test2"] = "test2"
assert new_model == expected_model, "This analysis invalidates everything"
def test_pipeline(storage_provider):
"""Load the schema and validate the pipeline.yml file against it."""
root = os.path.dirname(os.path.abspath(__file__))
pipeline = load_pipeline_yaml_file(os.path.join(root, "pipeline.yml"))
pipeline_configuration: PipelineConfiguration = {}
model = DictModel()
with open(os.path.join(root, "model.yml"), "rb") as model_file:
model.deserialize(model_file.read())
res = pipeline.get_artifact(
model=ReadOnlyModel(model),
artifact=pipeline.artifacts["ChildArtifact"],
requests=ObjectSet(MyKind.CHILD, {MyObjectID(MyKind.CHILD, "one")}),
pipeline_configuration={},
storage_provider=storage_provider,
)
assert res.objects() == ObjectSet(MyKind.CHILD, {MyObjectID(MyKind.CHILD, "one")})
res = pipeline.get_artifact(
model=ReadOnlyModel(model),
artifact=pipeline.artifacts["RootArtifact"],
requests=ObjectSet(MyKind.ROOT, {MyObjectID.root()}),
pipeline_configuration={},
storage_provider=storage_provider,
)
assert res.objects() == ObjectSet(MyKind.ROOT, {MyObjectID.root()})
new_model = pipeline.run_analysis(
model=ReadOnlyModel(model),
analysis_name="NullAnalysis",
requests=Requests(
{
ContainerDeclaration(
name="child_destination",
container_type=ChildDictContainer,
): ObjectSet(MyKind.CHILD, {MyObjectID(MyKind.CHILD, "one")})
}
),
analysis_configuration="",
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
assert isinstance(new_model, DictModel), "The analysis should return the same model type"
assert new_model == model, "NullAnalysis should not change the model"
new_model = pipeline.run_analysis(
model=ReadOnlyModel(model),
# An alias of PurgeAllAnalysis
analysis_name="blackhole",
requests=Requests(
{
ContainerDeclaration(
name="root_source",
container_type=RootDictContainer,
): ObjectSet(MyKind.ROOT, {MyObjectID.root()})
}
),
analysis_configuration="",
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
assert isinstance(new_model, DictModel), "The analysis should return the same model type"
assert len(new_model) == 0, "PurgeAllAnalysis should empty the model"
stored_model = storage_provider.get_model()
assert stored_model == new_model.serialize(), "The model should be stored"
def test_schedule_serdes(model):
storage_provider = InMemoryStorageProvider()
root1: ContainerDeclaration = ContainerDeclaration("root_source", RootDictContainer)
root2: ContainerDeclaration = ContainerDeclaration("root_destination", RootDictContainer)
child1: ContainerDeclaration = ContainerDeclaration("child_destination", ChildDictContainer)
child2: ContainerDeclaration = ContainerDeclaration("child_source", ChildDictContainer)
declarations = [root1, root2, child1, child2]
pipeline_configuration: PipelineConfiguration = {}
storage_provider.set_model(model.serialize())
begin_node = PipelineNode(SavePoint("begin", declarations))
up_node = PipelineNode(ToHigherKindPipe(), bindings=[root1, child1])
same_node = PipelineNode(SameKindPipe(), bindings=[child1, child2])
down_node = PipelineNode(ToLowerKindPipe(), bindings=[child2, root2])
end_node = PipelineNode(SavePoint("end", declarations))
begin_node.add_successor(up_node).add_successor(same_node).add_successor(
down_node
).add_successor(end_node)
pipeline: Pipeline = Pipeline(set(declarations), begin_node)
assert begin_node.savepoint_range == SavePointsRange(start=1, end=2)
assert up_node.savepoint_range == SavePointsRange(start=1, end=2)
assert same_node.savepoint_range == SavePointsRange(start=1, end=2)
assert down_node.savepoint_range == SavePointsRange(start=1, end=2)
assert end_node.savepoint_range == SavePointsRange(start=2, end=2)
root_obj = ObjectSet(MyKind.ROOT, {MyObjectID.root()})
# Force the savepoint to store the objects
container = RootDictContainer()
container.add_object(MyObjectID.root())
requests = Requests({root1: root_obj})
begin_node.run(
model=ReadOnlyModel(model),
containers={root1: container},
incoming=requests,
outgoing=requests,
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
schedule = pipeline.schedule(
model=ReadOnlyModel(model),
target_node=end_node,
requests=Requests({root2: root_obj}),
pipeline_configuration=pipeline_configuration,
storage_provider=storage_provider,
)
schedule_str = schedule.serialize()
pipeline.deserialize_schedule(schedule_str)